onetaskgraph_core/engine/copy.rs
1//! The copy verb: one item out of one source and into another, by the rules that make a
2//! second copy an update rather than a duplicate.
3//!
4//! Correspondence lives on the items and never in a table. A copied item carries
5//! [`GlobalId::ORIGIN_KEY`], whose value is the qualified id it was copied from, and the
6//! item it was copied from carries [`MetadataKey::COPIES_KEY`], an object naming its
7//! counterpart at each destination it was copied into. The rules below read exactly those —
8//! so nothing here is written down outside the plugin that owns the item, and the invariant
9//! this engine is built around is untouched.
10//!
11//! 0. **Follow the link.** An item whose `onetaskgraph.copies` names an item at the
12//! destination has that item read by id, and when that item's own origin names this one
13//! back it is the target: one read, and no scan. A link naming nothing there any more is
14//! refused as `stale-link` unless `--recreate` is given, for the reason a stale origin
15//! is; a link naming an item that no longer names this one back — somebody re-pointed it
16//! — is ignored, and the rules below decide.
17//! 1. **Follow the origin.** An item already carrying an origin whose source half is the
18//! destination names the destination item *directly*, and the copy updates it. This is
19//! the half that makes an edit's copy-back an update: the local file came from the
20//! remote item and knows which one.
21//! 2. **Search by origin.** Otherwise the destination is scanned, one page at a time, for
22//! an item whose origin is the id being copied. Found, the copy updates it; not found,
23//! the copy creates one carrying that origin.
24//!
25//! Which rule found the item decides what the copy records there. A copy that got its
26//! target from rule 1 is a copy-back — the destination is the *original*, and the item
27//! being copied is the one that came from it — so the destination keeps the origin it
28//! already holds, holding none included. Every other copy records the id it was copied
29//! from. See [`recorded`] for what stamping a copy-back's own id there costs.
30//!
31//! The same distinction decides the link. Once the whole copy has landed, every item whose
32//! target was found by any rule but rule 1 — or was created — has its link for this
33//! destination recorded or refreshed through its source's narrow metadata write, and that
34//! write joins the undo journal like every destination write. Rule 1 records nothing: that
35//! correspondence is already written down, on the destination item.
36//!
37//! A destination write is at the user's explicit request, names its destination, goes
38//! through that source's own write interface into that source's own store, and is never
39//! read back to answer a query. That is what makes it a write and not a cache.
40
41use std::collections::{BTreeMap, BTreeSet};
42
43use onetaskgraph_plugin_api::{
44 Capabilities, Cursor, DependencyEdge, DependencyEndpoint, DependencyKind, Direction, Document,
45 DocumentQuery, ItemKind, ItemWrite, Location, MetadataKey, MetadataMatch, Metering, NativeId,
46 Page, PageRequest, Project, ProjectQuery, Repository, SourceError, SourceName, StatusCategory,
47 Task, TaskQuery, TaskRef, TextFields, TextQuery,
48};
49use schemars::JsonSchema;
50use serde::{Deserialize, Serialize};
51use serde_json::Value;
52
53use crate::GlobalId;
54use crate::resolve::ResolvedSource;
55
56use super::delivery::{Delivered, targets};
57use super::fetch::{fits, unrepeated};
58use super::local::ProjectSelector;
59use super::narrow::holds_priority;
60use super::{
61 DocumentFilters, DocumentRequest, Engine, EngineError, Filters, LeftBehind, Paging, Qualified,
62 TaskRequest,
63};
64
65/// A request to copy work into one configured destination.
66#[derive(Debug, Clone)]
67pub struct CopyRequest {
68 /// The qualified items to copy, in the order they were named.
69 pub items: CopyItems,
70 /// What those ids name, and what comes with them.
71 pub scope: CopyScope,
72 /// The configured source to copy into — a source name, never a qualified id.
73 pub destination: SourceName,
74 /// How to re-establish a correspondence the two origin rules cannot find.
75 pub match_by: Option<MatchBy>,
76 /// Whether an origin naming nothing at the destination falls through to the search
77 /// rule instead of refusing.
78 pub recreate: bool,
79 /// Whether to perform every read and no write.
80 pub dry_run: bool,
81}
82
83/// The items one copy names: at least one, because a copy naming none is not a copy.
84///
85/// A newtype rather than a bare `Vec`, for the reason [`Repository`] is one: the empty
86/// list is not a copy of nothing, it is a caller mistake, and a type that can hold it
87/// leaves every reader to decide what it means — a report with no entries, an error, a
88/// silent success. None of those is better than not being able to say it.
89#[derive(Debug, Clone, PartialEq, Eq)]
90pub struct CopyItems(Vec<GlobalId>);
91
92impl CopyItems {
93 /// The items a caller named, or `None` when they named none.
94 #[must_use]
95 pub fn new(items: Vec<GlobalId>) -> Option<Self> {
96 (!items.is_empty()).then_some(Self(items))
97 }
98
99 /// The items, in the order they were named.
100 #[must_use]
101 pub fn as_slice(&self) -> &[GlobalId] {
102 &self.0
103 }
104}
105
106/// What the ids a copy names are, and what travels with them.
107///
108/// One value rather than a kind beside a flag, because three of the four combinations
109/// those two would make are real and the fourth — tasks, with the tasks of each also
110/// copied — means nothing. The member list is a variant for the same reason: it narrows a
111/// project copy that carries its tasks, and it has nothing to say to the other three.
112#[derive(Debug, Clone, PartialEq, Eq)]
113pub enum CopyScope {
114 /// The ids name tasks, and only those tasks are copied.
115 Tasks,
116 /// The ids name projects.
117 Projects {
118 /// Whether the tasks in each project are copied too.
119 tasks: bool,
120 },
121 /// The ids name projects, and of the tasks in them exactly these are copied.
122 ///
123 /// A member the list does not name is not read at the destination, not compared, not
124 /// written and not reported, and no walk for what the copy left behind runs. A named
125 /// member with no counterpart there is still created: the list narrows which members
126 /// are read and written, never which outcomes are possible.
127 ///
128 /// An edge from a copied item to a member the list does not name resolves to the
129 /// destination id that member's own [`GlobalId::ORIGIN_KEY`] records at the source,
130 /// without reading the destination for it. When that member records none, the copy is
131 /// refused before anything is written.
132 Members(CopyItems),
133 /// The ids name documents, and only those documents are copied.
134 ///
135 /// Nothing travels with a document: it takes part in no dependency graph, and it holds
136 /// nothing of its own the way a project holds tasks.
137 Documents,
138}
139
140/// The caller-named escape for a correspondence neither origin rule can find.
141///
142/// A person editing Markdown who deletes or corrupts the origin key leaves an item rule 1
143/// cannot use and rule 2 cannot find, and the next copy would create a second item. This
144/// is how that is re-established without hand-editing ids.
145#[derive(Debug, Clone, PartialEq, Eq)]
146pub enum MatchBy {
147 /// Match the item whose title is the same.
148 Title,
149 /// Match the item whose value at this metadata key is the same.
150 Metadata(String),
151}
152
153impl MatchBy {
154 /// The spelling a caller types, `title` or any metadata key.
155 #[must_use]
156 pub fn parse(key: &str) -> Self {
157 if key == "title" {
158 Self::Title
159 } else {
160 Self::Metadata(key.to_owned())
161 }
162 }
163}
164
165/// What a copy did, one entry per item.
166///
167/// The same per-item outcomes reach every consumer: the machine-readable output renders
168/// this, the rendered output renders this, and a Rust caller is handed it.
169#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
170pub struct CopyReport {
171 /// One entry per item the copy considered, in the order it considered them.
172 pub items: Vec<CopyOutcome>,
173 // Three flat scalars rather than one nested object, and a `!skip_serializing_if` beside
174 // every skip: both are load-bearing for the generated SDKs rather than matters of
175 // taste, and AGENTS.md's note on what a copied document's references are pointed at is
176 // where that reasoning lives.
177 /// Reference occurrences the copy rewrote to the destination's own location for the
178 /// record they name.
179 ///
180 /// A silent bound is indistinguishable from a bug, so the copy says what it did to the
181 /// references the documents it carried hold. This and the two below are totals over the
182 /// whole invocation rather than figures per document, and all three default to zero, so
183 /// a consumer written against the output before they existed is unaffected.
184 ///
185 /// **What these figures do not claim.** The referent set is a document's own project,
186 /// so a reference to a record in a *different* project is never recognised at all and
187 /// cannot appear in [`Self::references_unresolved`] either. These are the references
188 /// the copy recognised; they are not a census of every reference a document holds.
189 /// Noticing an out-of-scope reference would need exactly the unbounded destination walk
190 /// this design refuses.
191 #[serde(default, skip_serializing_if = "nothing_to_report")]
192 #[schemars(!skip_serializing_if)]
193 pub references_rewritten: u64,
194 /// Reference occurrences the copy recognised and left byte-for-byte as they were,
195 /// because the correspondence could not be established.
196 #[serde(default, skip_serializing_if = "nothing_to_report")]
197 #[schemars(!skip_serializing_if)]
198 pub references_unresolved: u64,
199 /// How many of [`Self::references_unresolved`] were left alone because the
200 /// correspondence was **ambiguous** rather than merely absent. A sub-count, never
201 /// larger than it.
202 ///
203 /// Split out because the two mean different things to a reader. A reference with no
204 /// counterpart is ordinary and expected under the bound above — the design working. An
205 /// ambiguous one says the destination holds duplicate records for one work item, or the
206 /// source reports one location for two records, and re-running the copy will never
207 /// clear it.
208 #[serde(default, skip_serializing_if = "nothing_to_report")]
209 #[schemars(!skip_serializing_if)]
210 // llmlint: ignore[invalid_states_unrepresentable] JSON Schema cannot express an
211 // inequality between two numbers, so a private constructor here would hold this in one
212 // consumer of three while both SDKs' generated models went on admitting it. What holds
213 // it is `substitute`: `Resolution` has no variant that counts an occurrence ambiguous
214 // without counting it unresolved.
215 pub references_ambiguous: u64,
216 /// `delivers` entries the copy rewrote to the destination's own id for a member of the
217 /// copied set, over the whole invocation. Every other entry arrives qualified and is not
218 /// counted. Left out when zero, as the three figures above are.
219 #[serde(default, skip_serializing_if = "nothing_to_report")]
220 #[schemars(!skip_serializing_if)]
221 pub delivers_rewritten: u64,
222 /// One entry per delivered task the copy kept in step with a task it landed, after the
223 /// whole copy was complete — see `task status set`, which reports the same entries. Left
224 /// out when there were none.
225 ///
226 /// A failed entry does not undo the copy: the tasks it landed stay landed, and the
227 /// command exits `4`.
228 #[serde(default, skip_serializing_if = "Vec::is_empty")]
229 #[schemars(!skip_serializing_if)]
230 pub delivered: Vec<Delivered>,
231 /// What this copy spent, summed over the sources in the command that meter their own
232 /// requests — and absent, never zero, when none of them does.
233 ///
234 /// Omitted from the wire when absent, like the three figures above, so a consumer
235 /// written against the output before it existed reads the same document.
236 #[serde(default, skip_serializing_if = "Option::is_none")]
237 pub spent: Option<Spent>,
238}
239
240/// Whether one of [`CopyReport`]'s reference figures has anything to say.
241///
242/// A copy that recognised no reference reports that by leaving the figure out rather than
243/// by writing a nought, so the machine output of a task or project copy is exactly what it
244/// was before these figures existed. The human rendering says it in words either way,
245/// because a reader there needs to be told the copy looked.
246fn nothing_to_report(figure: &u64) -> bool {
247 *figure == 0
248}
249
250/// What one command spent, summed over the sources in it that meter their own requests.
251///
252/// **Source-owned.** Every figure is what a source said it sent and spent while the command
253/// ran, read through [`TaskSource::metering`](onetaskgraph_plugin_api::TaskSource::metering)
254/// before the command and again after it. The engine adds the differences up by name and
255/// interprets none of them, which is why the budget and unit names are open vocabulary.
256#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
257pub struct Spent {
258 /// How many HTTP requests those sources sent for this command.
259 pub requests: u64,
260 /// What those requests spent, one entry per budget and unit, ordered by budget name.
261 pub budgets: Vec<BudgetSpent>,
262}
263
264/// What one command spent against one budget.
265#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
266pub struct BudgetSpent {
267 /// The budget, as the source names it — `graphql`, `rest`.
268 // llmlint: ignore[invalid_states_unrepresentable] The contract this report lands states `budget` and `unit` as strings in an open vocabulary the engine interprets none of, and both SDKs are generated from this type's schema; the one value such a vocabulary can refuse, the empty name, never reaches here, because `difference` below treats a source reporting one as not metering.
269 pub budget: String,
270 /// What it is metered in — `points`, `requests`.
271 // llmlint: ignore[invalid_states_unrepresentable] As `budget` above.
272 pub unit: String,
273 /// How much was spent against it, in that unit.
274 pub amount: u64,
275 /// Whether any part of `amount` was modelled by a source rather than reported by its
276 /// backend or counted, which makes `amount` a lower bound on what the backend charged
277 /// rather than a measurement of it.
278 pub lower_bound: bool,
279}
280
281/// One document's reference figures, before they are folded into the invocation's.
282#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
283struct Counted {
284 /// Occurrences rewritten.
285 rewritten: u64,
286 /// Occurrences recognised and left alone.
287 unresolved: u64,
288 /// How many of those were ambiguous.
289 ambiguous: u64,
290}
291
292impl Counted {
293 /// Fold one document's figures into the invocation's.
294 fn add(&mut self, other: Self) {
295 self.rewritten += other.rewritten;
296 self.unresolved += other.unresolved;
297 self.ambiguous += other.ambiguous;
298 }
299}
300
301/// What happened to one item.
302///
303/// `action` and `destination` are one value rather than two fields side by side: an
304/// updated item without a destination id, or an orphan without one, are states this type
305/// must not be able to say — the id *is* what those outcomes are about. The one outcome
306/// that legitimately has none is a dry run that would create, because nothing was
307/// created and there is no id to report.
308#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
309pub struct CopyOutcome {
310 /// The qualified id the item was read from.
311 pub source: GlobalId,
312 /// What happened to it, and where.
313 #[serde(flatten)]
314 pub action: CopyAction,
315}
316
317/// Which rule found the destination counterpart an item was updated at, or found already
318/// reading as it does.
319///
320/// Tried in this order, and the first that answers is the one reported: the link the item
321/// records, its own origin, a search of the destination for an item recording it as its
322/// origin, and the caller's `--match-by`. When none answers, the item is created and says
323/// so as [`NoCounterpart`] instead.
324#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
325#[serde(rename_all = "kebab-case")]
326pub enum CopyVia {
327 /// The item's `onetaskgraph.copies` entry for the destination named an item whose own
328 /// origin names this one back. Found by one read by id, and no scan.
329 Link,
330 /// The item's own `onetaskgraph.origin` named the destination item: a copy-back.
331 Origin,
332 /// A search of the destination for an item recording this one as its origin.
333 Scan,
334 /// A search of the destination by the caller's `--match-by`.
335 Match,
336}
337
338/// What a created item reports as its `via`: no rule found a counterpart.
339///
340/// One word, and a type of its own rather than a fifth [`CopyVia`], so an item updated at a
341/// counterpart cannot say none was found and a created one cannot name a rule that found
342/// one.
343#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
344#[serde(rename_all = "kebab-case")]
345pub enum NoCounterpart {
346 /// No counterpart was found, so one was created — or, in a dry run, would have been.
347 #[default]
348 Created,
349}
350
351/// What a copy did to the `onetaskgraph.copies` entry the copied item holds for the
352/// destination.
353#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
354#[serde(rename_all = "kebab-case")]
355pub enum CopyLink {
356 /// The entry was written, because the item held none for this destination or held one
357 /// naming another item.
358 Recorded,
359 /// Nothing was written: the entry already named this destination item, or the item was
360 /// found by its own origin, where the correspondence is already recorded on the other
361 /// side.
362 Unchanged,
363 /// The item's source cannot hold the entry — it has no write side, or it refuses the
364 /// narrow metadata write — so the copy landed without recording it.
365 Unrecorded,
366}
367
368impl CopyOutcome {
369 /// The qualified id this outcome landed on, when it landed on one.
370 #[must_use]
371 pub fn destination(&self) -> Option<&GlobalId> {
372 self.action.destination()
373 }
374}
375
376/// The four things a copy can do to one item, and the id each of them is about.
377#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
378#[serde(tag = "action", rename_all = "kebab-case")]
379pub enum CopyAction {
380 // llmlint: ignore[names_match_behavior] `created` is Contract D's serialized action
381 // for both a completed create and a dry run that would create; the optional destination
382 // distinguishes those cases, and renaming this public variant would break Rust callers.
383 /// The destination held no counterpart, so one was created.
384 Created {
385 /// The id it was created under, or `null` for a dry run that would have created
386 /// one — there is no id, because nothing was.
387 destination: Option<GlobalId>,
388 /// That no rule found a counterpart. The one word it can be, so a report written
389 /// before there was a `via` reads as saying it.
390 #[serde(default)]
391 via: NoCounterpart,
392 /// What the copy did to the link the copied item records for this destination at
393 /// `onetaskgraph.copies`; absent for a dry run, which writes nothing.
394 #[serde(default, skip_serializing_if = "Option::is_none")]
395 // llmlint: ignore[invalid_states_unrepresentable] Whether a copy is a dry run is a
396 // property of the request, which the report does not carry, and the contract this
397 // lands — shared with the follow-up tooling built against it — says a dry run leaves
398 // `link` out rather than naming a sixth word for it. What holds it is `Engine::link`:
399 // the one writer of this field, run for every item a copy that writes has landed and
400 // for none of a dry run's, which `a_dry_run_says_which_rule_would_answer_and_records_no_link`
401 // and every `link` assertion in tests/copy_link.rs drive.
402 link: Option<CopyLink>,
403 },
404 /// The destination held a counterpart and it now reads as the source does.
405 Updated {
406 /// The item that was updated.
407 destination: GlobalId,
408 /// Which rule found it.
409 via: CopyVia,
410 /// What the copy did to the copied item's link; absent for a dry run.
411 #[serde(default, skip_serializing_if = "Option::is_none")]
412 // llmlint: ignore[invalid_states_unrepresentable] As `link` on `Created` above.
413 link: Option<CopyLink>,
414 },
415 /// The destination held a counterpart that already read that way; nothing was written.
416 Unchanged {
417 /// The item that already said it.
418 destination: GlobalId,
419 /// Which rule found it.
420 via: CopyVia,
421 /// What the copy did to the copied item's link; absent for a dry run.
422 #[serde(default, skip_serializing_if = "Option::is_none")]
423 // llmlint: ignore[invalid_states_unrepresentable] As `link` on `Created` above.
424 link: Option<CopyLink>,
425 },
426 /// The destination holds a counterpart the source no longer does. A copy never
427 /// deletes, so it was left exactly as it is — and no rule was asked where it goes, and
428 /// no link was touched, because it was not copied.
429 Orphaned {
430 /// The item that was left alone.
431 destination: GlobalId,
432 },
433}
434
435impl CopyAction {
436 /// The qualified id this action is about, when there is one.
437 #[must_use]
438 pub fn destination(&self) -> Option<&GlobalId> {
439 match self {
440 Self::Created { destination, .. } => destination.as_ref(),
441 Self::Updated { destination, .. }
442 | Self::Unchanged { destination, .. }
443 | Self::Orphaned { destination } => Some(destination),
444 }
445 }
446
447 /// Which rule found the counterpart, or `None` when none did — the item was created —
448 /// and for an orphan, which no rule was asked about.
449 #[must_use]
450 pub fn found_by(&self) -> Option<CopyVia> {
451 match self {
452 Self::Updated { via, .. } | Self::Unchanged { via, .. } => Some(*via),
453 Self::Created { .. } | Self::Orphaned { .. } => None,
454 }
455 }
456
457 /// What became of the copied item's link, or `None` for a dry run and for an orphan.
458 #[must_use]
459 pub fn link(&self) -> Option<CopyLink> {
460 match self {
461 Self::Created { link, .. }
462 | Self::Updated { link, .. }
463 | Self::Unchanged { link, .. } => *link,
464 Self::Orphaned { .. } => None,
465 }
466 }
467
468 /// The link slot of an item this copy landed, `None` for an orphan.
469 fn link_slot(&mut self) -> Option<&mut Option<CopyLink>> {
470 match self {
471 Self::Created { link, .. }
472 | Self::Updated { link, .. }
473 | Self::Unchanged { link, .. } => Some(link),
474 Self::Orphaned { .. } => None,
475 }
476 }
477
478 /// The word this action serializes as, taken from its own `Serialize`.
479 ///
480 /// Read back off the wire form rather than written out again in a `match`, for the
481 /// reason `render::wire` gives: a second spelling of `unchanged` would be a second
482 /// place for it to drift from the one a caller reads.
483 #[must_use]
484 pub fn name(&self) -> String {
485 serde_json::to_value(self).expect("a contract enum serialises")["action"]
486 .as_str()
487 .expect("an internally tagged enum carries its tag")
488 .to_owned()
489 }
490}
491
492/// What one item points at, resolved as far as the copy has got: its forward edges and the
493/// tasks it delivers.
494struct Pointing<'a> {
495 /// Its forward edges, `None` where the far end is a member not landed yet.
496 edges: &'a [Option<DependencyEdge>],
497 /// The tasks it delivers that can already be named at the destination.
498 delivers: &'a [TaskRef],
499}
500
501/// Where one item is going at the destination.
502enum Target {
503 /// Update the destination item with this id, reached by the rule named.
504 Update {
505 /// The destination item this copy updates.
506 id: NativeId,
507 /// Which rule found it.
508 found: Found,
509 },
510 /// Create one.
511 Create,
512}
513
514/// Which of the rules above found the destination item a copy is updating.
515///
516/// Every one is the same instruction — update that item — and rule 1 is a different answer
517/// about the origin, which is why the distinction is carried this far rather than dropped
518/// where it is made. See [`recorded`]. It is also what the report's [`CopyVia`] says and
519/// what decides whether the copied item's link is written.
520#[derive(Clone, Copy, PartialEq, Eq)]
521enum Found {
522 /// The item's own link named it, and it names the item back.
523 Link,
524 /// Rule 1: the item being copied already named it, so this copy is a copy-back.
525 Origin,
526 /// Rule 2: the destination was searched for an item recording this one as its origin.
527 Scan,
528 /// The caller's matching escape.
529 Match,
530}
531
532impl Found {
533 /// The word the report says it in.
534 fn via(self) -> CopyVia {
535 match self {
536 Self::Link => CopyVia::Link,
537 Self::Origin => CopyVia::Origin,
538 Self::Scan => CopyVia::Scan,
539 Self::Match => CopyVia::Match,
540 }
541 }
542}
543
544impl Target {
545 /// Which rule found this target, or `None` for one this copy creates.
546 fn found(&self) -> Option<Found> {
547 match self {
548 Self::Update { found, .. } => Some(*found),
549 Self::Create => None,
550 }
551 }
552}
553
554/// What a scan of the destination is looking for.
555enum Wanted {
556 /// An item recording this qualified id as its origin.
557 Origin(String),
558 /// An item whose title is this.
559 Title(String),
560 /// An item holding this value at this metadata key.
561 Metadata(String, Value),
562}
563
564impl Wanted {
565 /// Whether one destination item is the one being looked for.
566 fn found(&self, title: &str, metadata: &BTreeMap<String, Value>) -> bool {
567 match self {
568 Self::Origin(id) => {
569 metadata.get(GlobalId::ORIGIN_KEY) == Some(&Value::String(id.clone()))
570 }
571 Self::Title(wanted) => title == wanted,
572 Self::Metadata(key, value) => metadata.get(key) == Some(value),
573 }
574 }
575
576 /// The task query that asks `capabilities`' source for this item rather than for all of
577 /// them, when that source applies the predicate itself.
578 ///
579 /// Only ever a narrowing the source declared it applies, so a source that cannot is
580 /// asked for everything exactly as before; and never the last word, because
581 /// [`found`](Self::found) still confirms each row. A text search is wider than an equal
582 /// title, and a source may answer a match it cannot yet see as absent — which is the
583 /// same answer the whole-store walk gave for an item its own listing was behind on.
584 ///
585 /// A title or a value with no letter or digit is never pushed down: a source whose search
586 /// indexes words cannot answer it and may refuse it, and the walk it replaces answered it
587 /// already, so it is asked for everything as it always was.
588 fn task_query(&self, capabilities: &Capabilities) -> TaskQuery {
589 match self {
590 Self::Origin(id) if capabilities.filter_by_origin.is_native() => TaskQuery {
591 origin: Some(id.clone()),
592 ..TaskQuery::default()
593 },
594 Self::Title(title) if capabilities.search_title.is_native() && has_words(title) => {
595 TaskQuery {
596 text: Some(TextQuery {
597 terms: title.clone(),
598 fields: TextFields::Title,
599 }),
600 ..TaskQuery::default()
601 }
602 }
603 Self::Metadata(key, Value::String(value))
604 if capabilities.filter_by_metadata.is_native() && has_words(value) =>
605 {
606 MetadataMatch::new(key.clone(), Vec::new(), value.clone())
607 .map(|wanted| TaskQuery {
608 metadata: vec![wanted],
609 ..TaskQuery::default()
610 })
611 .unwrap_or_default()
612 }
613 _ => TaskQuery::default(),
614 }
615 }
616}
617
618/// Whether `terms` holds a letter or a digit — a word a search that indexes words can find.
619fn has_words(terms: &str) -> bool {
620 terms.chars().any(char::is_alphanumeric)
621}
622
623/// What the destination held before this copy touched one item.
624///
625/// Read once, in [`Engine::land`], and used three times over: to decide whether the write
626/// would change anything, to repair the item's edges once the rest of the copy has landed,
627/// and — if the copy cannot finish — to put the item back exactly as it was.
628#[derive(Clone)]
629struct Prior {
630 /// The item as the destination held it.
631 item: Item,
632 /// Its forward edges there.
633 edges: Vec<DependencyEdge>,
634}
635
636/// One item that landed with an edge whose far end was not written yet.
637///
638/// Held until every item of the whole copy has landed, because the far end may be in
639/// another project of the same command: a copy of two projects at once is one copied set,
640/// not two, and an edge across them is remapped rather than written as a foreign id.
641struct Deferred {
642 /// The item, as it was read and resolved.
643 item: Planned,
644 /// The destination project it was filed under.
645 filed: Option<NativeId>,
646 /// Where it landed.
647 destination: NativeId,
648 /// What the destination held there before, when it held anything.
649 prior: Option<Prior>,
650}
651
652/// What one item's undo has to do to put the destination back.
653enum Undo {
654 /// The copy created it, so undoing means removing it.
655 Created {
656 /// Which write interface removes it.
657 kind: Level,
658 /// The destination id it was created under.
659 id: NativeId,
660 },
661 /// The copy overwrote something, so undoing means writing that something back.
662 ///
663 /// No `kind` beside the id, unlike the variant above: what was there says which of the
664 /// two write interfaces takes it back, and a second spelling of that could disagree
665 /// with it.
666 Updated {
667 /// The destination id that was overwritten.
668 id: NativeId,
669 /// What was there before.
670 prior: Prior,
671 },
672}
673
674impl Undo {
675 /// The destination id this entry is about.
676 fn id(&self) -> &NativeId {
677 match self {
678 Self::Created { id, .. } | Self::Updated { id, .. } => id,
679 }
680 }
681
682 /// Which of the destination's three write interfaces this entry belongs to.
683 ///
684 /// An id alone does not identify a destination item: nothing stops a destination
685 /// numbering its tasks and its projects in one namespace, and a local-Markdown store
686 /// filing `alpha.md` under both is the ordinary case rather than the contrived one.
687 /// So this pairs with `id` wherever one entry has to be told from another.
688 fn kind(&self) -> Level {
689 match self {
690 Self::Created { kind, .. } => *kind,
691 // Read off what was there, for the reason the variant carries no `kind` of
692 // its own: two spellings of one fact can disagree, and this one cannot.
693 Self::Updated { prior, .. } => prior.item.level(),
694 }
695 }
696}
697
698/// Everything one copy has written, in the order it wrote it, so a copy that cannot finish
699/// can undo its own writes.
700///
701/// This is not state the engine keeps: it lives for the length of one `copy` call and is
702/// dropped with it, so the invariant that nothing of a user's work is written down outside
703/// the plugin that owns it is untouched.
704#[derive(Default)]
705struct Journal {
706 /// One entry per destination item this copy first touched, in that order.
707 entries: Vec<Undo>,
708 /// One entry per copied item whose `onetaskgraph.copies` link this copy wrote, in that
709 /// order.
710 ///
711 /// Kept apart from `entries` because these are the *sources'* items rather than the
712 /// destination's, and an id alone does not say which store it is in: a copy inside one
713 /// source can write a link on an item and land on another item with the same kind.
714 links: Vec<Unlink>,
715}
716
717/// How to take back one link this copy wrote on a copied item.
718struct Unlink {
719 /// The item, at its own source.
720 item: GlobalId,
721 /// Which of that source's interfaces holds it.
722 level: Level,
723 /// What it held at `onetaskgraph.copies` before, or `None` when it held nothing there.
724 before: Option<Value>,
725}
726
727impl Journal {
728 /// Record what has to happen to put one destination item back.
729 ///
730 /// The *first* entry for an id is the one that matters and later ones are dropped: an
731 /// item written twice — once as it lands, once when its edges are repaired — was only
732 /// ever one thing before this copy started, and that is what undoing it restores.
733 fn record(&mut self, entry: Undo) {
734 if self
735 .entries
736 .iter()
737 .any(|held| held.kind() == entry.kind() && held.id() == entry.id())
738 {
739 return;
740 }
741 self.entries.push(entry);
742 }
743}
744
745/// What one copy invocation carries from the item it lands to the next.
746///
747/// Like the [`Journal`] beside it, this lives for the length of one `copy` call and is
748/// dropped with it: nothing here is state the engine keeps, and nothing in it is read back
749/// to answer a later query.
750#[derive(Default)]
751struct Running {
752 /// Every item an edge of this command resolves to a destination id rather than writing as
753 /// it was read: the ids named, what travels with them, and — for a member copy — each
754 /// member the copy does not carry whose own recorded origin already names its
755 /// destination item. That last kind is resolvable without being copied.
756 resolvable: Vec<GlobalId>,
757 /// The destination item each of those corresponds to, as soon as it is known — whether
758 /// this command found it, landed it, or read it off a member's recorded origin.
759 ///
760 /// Keyed by the qualified id's own rendering, which is what a recorded origin holds
761 /// anyway — making `GlobalId` orderable for a local map would put an ordering on a
762 /// contract type for a reason no caller of it has.
763 counterparts: BTreeMap<String, NativeId>,
764 /// The items held back for [`Engine::repair`], across the whole request.
765 deferred: Vec<Deferred>,
766 /// The reference figures, one total for the whole invocation rather than one per
767 /// document.
768 references: Counted,
769 /// The destination project each source project corresponds to, once this command has
770 /// looked for it — so a second task filed under the same project does not walk the
771 /// destination for it again.
772 filings: BTreeMap<String, Option<NativeId>>,
773 /// `delivers` entries naming a member of the copied set, over the whole invocation.
774 delivers_rewritten: u64,
775 /// Every task this copy landed, for the relation rule once the whole copy is complete.
776 landed: Vec<LandedTask>,
777 /// Every item this copy landed, in the order it landed them, for the link each records
778 /// once the whole copy has landed.
779 linking: Vec<Linking>,
780}
781
782/// One item a copy landed, and what settling its `onetaskgraph.copies` link needs.
783struct Linking {
784 /// The item, at its own source.
785 item: GlobalId,
786 /// Which of that source's interfaces holds it.
787 level: Level,
788 /// Which rule found where it landed, or `None` when this copy created it.
789 found: Option<Found>,
790 /// Where it landed.
791 destination: GlobalId,
792 /// What it held at `onetaskgraph.copies` when it was read, if anything.
793 held: Option<Value>,
794}
795
796/// One task a copy landed, and what the relation rule reads of it once the copy is complete.
797struct LandedTask {
798 /// Where it landed.
799 destination: GlobalId,
800 /// The source it was read from, which a bare `delivers` entry names a task of.
801 origin: SourceName,
802 /// Its `delivers`, as its source reported them.
803 delivers: Vec<TaskRef>,
804 /// The `delivers` the destination held there before this copy, qualified.
805 before: Vec<GlobalId>,
806 /// Its status category, as the copy wrote it.
807 category: StatusCategory,
808}
809
810/// A task a copy landed, with the tasks it delivers now and delivered before, qualified.
811struct Deliverer {
812 destination: GlobalId,
813 category: StatusCategory,
814 now: Vec<GlobalId>,
815 before: Vec<GlobalId>,
816}
817
818/// One item, read and resolved, on its way into the destination.
819struct Planned {
820 /// Where it came from.
821 source: GlobalId,
822 /// The item as its source reported it.
823 item: Item,
824 /// Its forward edges, as its source reported them.
825 edges: Vec<DependencyEdge>,
826 /// Where it is going.
827 target: Target,
828 /// What the destination held at that target, read once where the target was found.
829 ///
830 /// The read that decided an item exists is the read of what it holds, so the two are
831 /// one round trip rather than two against a hosted destination.
832 held: Option<Prior>,
833}
834
835/// A task, a project or a document, so the copy path is written once.
836#[derive(Clone)]
837enum Item {
838 /// A task.
839 Task(Box<Task>),
840 /// A project.
841 Project(Box<Project>),
842 /// A document.
843 Document(Box<Document>),
844}
845
846impl Item {
847 fn id(&self) -> &NativeId {
848 match self {
849 Self::Task(task) => &task.id,
850 Self::Project(project) => &project.id,
851 Self::Document(document) => &document.id,
852 }
853 }
854
855 fn level(&self) -> Level {
856 match self {
857 Self::Task(_) => Level::Task,
858 Self::Project(_) => Level::Project,
859 Self::Document(_) => Level::Document,
860 }
861 }
862}
863
864/// Which of a destination's three read-and-write interfaces one item belongs to.
865///
866/// Deliberately not [`ItemKind`]: that enum names what a *dependency endpoint* points at,
867/// and the contract gives it no document variant because nothing may point at a document.
868/// This one names which pair of methods reads and writes an item, which is a different
869/// question with a third answer.
870///
871/// Ordered so it can key a map of what a destination holds, per interface: an id alone
872/// does not identify a destination item, for the reason [`Undo::kind`] records.
873#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
874enum Level {
875 /// `get_task`, `write_task`, `delete_task`.
876 Task,
877 /// `get_project`, `write_project`, `delete_project`.
878 Project,
879 /// `get_document`, `write_document`, `delete_document`.
880 Document,
881}
882
883/// One record filed under a document's own project, and the location string its source
884/// reports for it.
885///
886/// A reference is a **literal occurrence, in a document's content, of the exact location
887/// string the source reports for a related record** — the `String` inside
888/// [`Location::Path`] or [`Location::Url`]. Both ends of a rewrite come from the plugins'
889/// own reported [`Location`]: nothing here composes an address out of a name, an id or a
890/// root, because a source that reports a canonical absolute path and a source that reports
891/// an issue link are the two things the contract lets this ask about.
892#[derive(Clone)]
893struct Referent {
894 /// Its qualified id at the source.
895 id: GlobalId,
896 /// The origin it records in its own metadata, when it records one.
897 origin: Option<GlobalId>,
898 /// Which of the destination's three interfaces its counterpart would be read from.
899 level: Level,
900 /// The non-empty location string its source reports for it.
901 location: String,
902}
903
904impl Referent {
905 /// The two keys a destination record's own recorded origin is matched against.
906 ///
907 /// **A destination record is this referent's counterpart when its
908 /// [`GlobalId::ORIGIN_KEY`] equals either this referent's own qualified source id, or
909 /// the origin this referent itself records.** In plainer terms: when the destination
910 /// record was copied directly from this referent, or directly from the one predecessor
911 /// this referent itself records.
912 ///
913 /// That reach is exactly one recorded hop of ancestry on each side, and nothing more:
914 ///
915 /// - **A single hop resolves.** The destination record came straight from the referent.
916 /// - **A one-level fan-out resolves.** A record copied from one store into two, with
917 /// the document arriving by one route and naming records that arrived by the other,
918 /// so both sides trace to one common predecessor. That is what the second key buys,
919 /// and it is the only thing it buys.
920 /// - **A chain of two or more hops does not resolve**, on either side, and that is
921 /// permanent rather than pending: only one origin is ever recorded and every hop
922 /// overwrites it. [`carried`] removes the key from the incoming metadata outright and
923 /// [`recorded`] writes the id at the *immediate* source on every path but the
924 /// copy-back, so A → B → C leaves C keyed by B and A's id gone.
925 ///
926 /// The second key costs no read — the referent's metadata is already in hand. Chasing
927 /// the chain further would need the intermediate stores configured and reachable, which
928 /// would put a third party's availability inside a copy; a durable lineage id would
929 /// identify only records written after it landed, so it would resolve nothing already
930 /// on a destination.
931 fn keys(&self) -> Vec<String> {
932 let mut keys = vec![self.id.to_string()];
933 if let Some(origin) = &self.origin {
934 let recorded = origin.to_string();
935 if recorded != keys[0] {
936 keys.push(recorded);
937 }
938 }
939 keys
940 }
941}
942
943/// What every whole-reference occurrence of one location string becomes.
944///
945/// Three variants rather than a rewrite beside a flag, because the two ways of leaving an
946/// occurrence alone are what the two figures a copy reports are *about*: one is the design
947/// working and the other says the destination or the source holds something a re-run will
948/// never fix. A bool would let a reader of this type read them as the same outcome.
949enum Resolution {
950 /// The destination's own location string for the counterpart.
951 Rewrite(String),
952 /// The destination holds no counterpart, or holds one it reports no location for, so
953 /// the occurrence is left byte-for-byte. Ordinary and expected under the bound this
954 /// design works to.
955 NoCounterpart,
956 /// The correspondence could not be established **confidently** — more than one
957 /// destination record matches the referent's two keys, or two referents report this one
958 /// location string. The occurrence is left byte-for-byte, no record is chosen, and a
959 /// re-run will never clear it.
960 Ambiguous,
961}
962
963/// Every destination record that records an origin, read **once per copy invocation**.
964///
965/// Several documents in one project is the ordinary case, and a walk per document would
966/// multiply reads against a rate limiter for no gain — so this is built once, for the
967/// levels the invocation's own documents really name, and every document and every
968/// referent of that invocation is answered out of it. A copy whose documents hold no
969/// candidate reference builds none at all.
970///
971/// **This is a stricter discipline than [`Engine::scan`], deliberately.** That lookup takes
972/// the first hit and stops, and every consumer of the copy already depends on it doing so;
973/// it chooses the copy's own target, where a caller named the item. This one edits the
974/// content of somebody's document, where a wrong answer is silent corruption of prose a
975/// person will act on — so where more than one record matches, it chooses none. The two
976/// lookups answer different questions and are meant to disagree on a destination holding
977/// duplicates.
978#[derive(Default)]
979struct Counterparts {
980 /// The records at one interface recording one origin.
981 by_origin: BTreeMap<(Level, String), Vec<Held>>,
982}
983
984/// One record a destination walk found: where it is there, and the location string that
985/// destination reports for it — `None` when it reports none.
986type Held = (NativeId, Option<String>);
987
988impl Counterparts {
989 /// Record one destination item, when it records an origin at all.
990 fn note(
991 &mut self,
992 level: Level,
993 id: &NativeId,
994 location: Option<&Location>,
995 metadata: &BTreeMap<String, Value>,
996 ) {
997 let Some(origin) = origin_of(metadata) else {
998 return;
999 };
1000 self.by_origin
1001 .entry((level, origin.to_string()))
1002 .or_default()
1003 .push((id.clone(), located(location)));
1004 }
1005
1006 /// What one referent's occurrences become, by the two-key rule.
1007 ///
1008 /// Where the correspondence cannot be established **confidently**, the text is left
1009 /// exactly as it is and no record is chosen — see the note on this type.
1010 fn resolve(&self, referent: &Referent) -> Resolution {
1011 let mut candidates: Vec<&Held> = Vec::new();
1012 for key in referent.keys() {
1013 for record in self
1014 .by_origin
1015 .get(&(referent.level, key))
1016 .into_iter()
1017 .flatten()
1018 {
1019 // One destination record matching both keys is one record, not two.
1020 if !candidates.iter().any(|held| held.0 == record.0) {
1021 candidates.push(record);
1022 }
1023 }
1024 }
1025 match candidates.as_slice() {
1026 [] => Resolution::NoCounterpart,
1027 // A counterpart the destination reports no location for names nowhere a reader
1028 // could go, so the source's own string is left standing rather than removed.
1029 [(_, location)] => location
1030 .clone()
1031 .map_or(Resolution::NoCounterpart, Resolution::Rewrite),
1032 _ => Resolution::Ambiguous,
1033 }
1034 }
1035}
1036
1037impl Engine {
1038 /// Copy every item a request names into one configured destination.
1039 ///
1040 /// This is the whole of the verb, and the command line drives exactly this: a copy a
1041 /// Rust caller makes and a copy typed at a shell are the same call, so the two cannot
1042 /// answer the same copy differently.
1043 ///
1044 /// # Errors
1045 ///
1046 /// Returns [`EngineError`] when the destination is not configured, cannot be built,
1047 /// cannot be written, or — for a document copy — declares it has no documents; when an
1048 /// id names nothing; when an origin names an item the
1049 /// destination no longer holds and `--recreate` was not given; and when the
1050 /// destination refuses the write — including a field or a metadata key it cannot
1051 /// carry, which it names rather than dropping.
1052 pub async fn copy(&self, request: &CopyRequest) -> Result<CopyReport, EngineError> {
1053 let destination = self.writable(&request.destination)?;
1054 // Before anything is read, and from the declaration rather than from a failed
1055 // write: a destination that says it has no documents has nowhere to put one.
1056 if request.scope == CopyScope::Documents {
1057 documentary(destination)?;
1058 }
1059 // The sources this command names, each read once before and once after, so what it
1060 // spent is the difference between two of each one's own running totals.
1061 let mut metered = vec![destination];
1062 for id in request.items.as_slice() {
1063 if let Some(source) = self.ready().find(|source| source.name() == &id.source)
1064 && !metered.iter().any(|held| held.name() == source.name())
1065 {
1066 metered.push(source);
1067 }
1068 }
1069 let before = readings(&metered).await;
1070 let mut journal = Journal::default();
1071 match self.copy_all(destination, request, &mut journal).await {
1072 Ok((mut report, deliverers)) => {
1073 // After the whole copy is complete and outside the journal: what the relation
1074 // rule writes is the delivered tasks' own, and a failure there is reported
1075 // against that task rather than undoing a copy that landed.
1076 for deliverer in &deliverers {
1077 report.delivered.extend(
1078 self.deliver(
1079 &deliverer.destination,
1080 deliverer.category,
1081 &deliverer.now,
1082 &deliverer.before,
1083 )
1084 .await,
1085 );
1086 }
1087 report.spent = spent_between(&before, &readings(&metered).await);
1088 Ok(report)
1089 }
1090 Err(error) => Err(self.undo(destination, journal, error).await),
1091 }
1092 }
1093
1094 /// The copy itself, with everything it writes recorded so a failure can be undone.
1095 ///
1096 /// The ids named together are **one** copied set, and that is what makes an edge
1097 /// between any two of them a real edge at the destination: a copy of two projects at
1098 /// once knows that a task in the first depends on a task in the second, and a task
1099 /// knows that the project it belongs to is being created beside it. Copying them one
1100 /// at a time could not, and wrote the far end as the id it had at its *source* — a
1101 /// dangling reference to somewhere the destination has never heard of.
1102 async fn copy_all(
1103 &self,
1104 destination: &ResolvedSource,
1105 request: &CopyRequest,
1106 journal: &mut Journal,
1107 ) -> Result<(CopyReport, Vec<Deliverer>), EngineError> {
1108 let mut running = Running::default();
1109 // The whole copied set, established before anything is written. For a project
1110 // copy that means reading every named project's membership first: the set is the
1111 // whole request rather than one project of it.
1112 let items = match &request.scope {
1113 CopyScope::Tasks | CopyScope::Documents => {
1114 let kind = if request.scope == CopyScope::Documents {
1115 Level::Document
1116 } else {
1117 Level::Task
1118 };
1119 running
1120 .resolvable
1121 .extend(request.items.as_slice().iter().cloned());
1122 let mut planned = Vec::new();
1123 for id in request.items.as_slice() {
1124 planned.push(self.plan(destination, request, kind, id).await?);
1125 }
1126 self.copy_items(destination, request, planned, None, &mut running, journal)
1127 .await?
1128 }
1129 CopyScope::Projects { tasks } => {
1130 let mut projects = Vec::new();
1131 for id in request.items.as_slice() {
1132 let members = if *tasks {
1133 self.project_members(id).await?
1134 } else {
1135 Vec::new()
1136 };
1137 running.resolvable.push(id.clone());
1138 running.resolvable.extend(members.iter().cloned());
1139 projects.push((id.clone(), members));
1140 }
1141 self.copy_projects(destination, request, &projects, &[], &mut running, journal)
1142 .await?
1143 }
1144 CopyScope::Members(named) => {
1145 let (projects, unrecorded) = self
1146 .named_members(destination, request.items.as_slice(), named, &mut running)
1147 .await?;
1148 self.copy_projects(
1149 destination,
1150 request,
1151 &projects,
1152 &unrecorded,
1153 &mut running,
1154 journal,
1155 )
1156 .await?
1157 }
1158 };
1159 let references = running.references;
1160 let delivers_rewritten = running.delivers_rewritten;
1161 // Every member's destination id is known once the first pass has landed it, so the
1162 // tasks each landed task delivers are settled here rather than after the repair.
1163 let deliverers: Vec<Deliverer> = running
1164 .landed
1165 .iter()
1166 .map(|task| Deliverer {
1167 destination: task.destination.clone(),
1168 category: task.category,
1169 now: targets(
1170 &resolved_entries(&mapped_delivers(
1171 &task.delivers,
1172 &task.origin,
1173 destination,
1174 &running.resolvable,
1175 &running.counterparts,
1176 )),
1177 destination.name(),
1178 ),
1179 before: task.before.clone(),
1180 })
1181 .collect();
1182 let linking = std::mem::take(&mut running.linking);
1183 self.repair(destination, request, running, journal).await?;
1184 let mut items = items;
1185 self.link(destination, linking, &mut items, journal).await?;
1186 Ok((
1187 CopyReport {
1188 items,
1189 references_rewritten: references.rewritten,
1190 references_unresolved: references.unresolved,
1191 references_ambiguous: references.ambiguous,
1192 delivers_rewritten,
1193 // Filled in by `copy`, once the copy is complete.
1194 delivered: Vec::new(),
1195 // Filled in by `copy`, which is the one place both readings are taken.
1196 spent: None,
1197 },
1198 deliverers,
1199 ))
1200 }
1201
1202 /// Write every deferred item again, now that every destination id is known.
1203 ///
1204 /// This is the second half of the two passes an edge between two items of one copy
1205 /// needs: the far end's destination id does not exist until it has been created, so
1206 /// the item that points at it lands first without that edge and is completed here.
1207 /// It runs once for the whole request rather than once per project, because a far end
1208 /// may be in a project this copy has not reached yet.
1209 async fn repair(
1210 &self,
1211 destination: &ResolvedSource,
1212 request: &CopyRequest,
1213 running: Running,
1214 journal: &mut Journal,
1215 ) -> Result<(), EngineError> {
1216 if request.dry_run {
1217 return Ok(());
1218 }
1219 let Running {
1220 resolvable,
1221 counterparts,
1222 deferred,
1223 ..
1224 } = running;
1225 for entry in deferred {
1226 let edges = mapped_edges(
1227 &entry.item.edges,
1228 &entry.item.source.source,
1229 destination,
1230 &resolvable,
1231 &counterparts,
1232 );
1233 let delivers = delivers_of(&entry.item, destination, &resolvable, &counterparts);
1234 self.write(
1235 destination,
1236 &entry.item,
1237 Some(entry.destination),
1238 entry.filed,
1239 &resolved(&edges),
1240 &resolved_entries(&delivers),
1241 entry.prior,
1242 journal,
1243 )
1244 .await?;
1245 }
1246 Ok(())
1247 }
1248
1249 /// Put the destination back the way this copy found it, then report why it failed.
1250 ///
1251 /// Undone in reverse, and an item this copy created is removed rather than restored —
1252 /// the entry recording what it looked like a moment after creation is not a state
1253 /// anybody asked for. When the destination cannot take one of them back, the refusal
1254 /// says so and names what is still there, because a user told "the copy failed" about
1255 /// a destination that is not as they left it will copy again over a tree nobody
1256 /// described.
1257 async fn undo(
1258 &self,
1259 destination: &ResolvedSource,
1260 journal: Journal,
1261 error: EngineError,
1262 ) -> EngineError {
1263 let created: Vec<(Level, &NativeId)> = journal
1264 .entries
1265 .iter()
1266 .filter_map(|entry| match entry {
1267 Undo::Created { kind, id } => Some((*kind, id)),
1268 Undo::Updated { .. } => None,
1269 })
1270 .collect();
1271 // The ids and the refusal are one value rather than two, because they are one
1272 // fact: an item is only left behind because the destination refused to take it
1273 // back, so the first refusal carries the first id and neither half can be
1274 // recorded without the other.
1275 let mut unrestored: Option<(LeftBehind, SourceError)> = None;
1276 // The links first, because they were written last: once the whole copy had landed.
1277 for entry in journal.links.iter().rev() {
1278 if let Err(problem) = self.unlink(entry).await {
1279 match &mut unrestored {
1280 Some((left_behind, _)) => left_behind.push(entry.item.clone()),
1281 None => unrestored = Some((LeftBehind::new(entry.item.clone()), problem)),
1282 }
1283 }
1284 }
1285 for entry in journal.entries.iter().rev() {
1286 let outcome = match entry {
1287 Undo::Created { kind, id } => remove(destination, *kind, id).await,
1288 Undo::Updated { id, prior, .. } if !created.contains(&(prior.item.level(), id)) => {
1289 restore(destination, id, prior).await
1290 }
1291 Undo::Updated { .. } => Ok(()),
1292 };
1293 if let Err(problem) = outcome {
1294 let id = GlobalId::new(destination.name().clone(), entry.id().clone());
1295 match &mut unrestored {
1296 Some((left_behind, _)) => left_behind.push(id),
1297 None => unrestored = Some((LeftBehind::new(id), problem)),
1298 }
1299 }
1300 }
1301 match unrestored {
1302 None => error,
1303 Some((left_behind, refusal)) => EngineError::CopyNotUndone {
1304 error: Box::new(error),
1305 left_behind,
1306 refusal,
1307 },
1308 }
1309 }
1310
1311 /// Record, on every item this copy landed, where it landed — once the whole copy has
1312 /// landed, so a copy that fails before this point has written nothing on any source.
1313 ///
1314 /// Each outcome is told what happened to its item's link. The items landed are exactly
1315 /// the outcomes that are not orphans, in the order they landed; a dry run records
1316 /// nothing, so this is never asked about one.
1317 async fn link(
1318 &self,
1319 destination: &ResolvedSource,
1320 linking: Vec<Linking>,
1321 items: &mut [CopyOutcome],
1322 journal: &mut Journal,
1323 ) -> Result<(), EngineError> {
1324 let mut landed = items
1325 .iter_mut()
1326 .filter(|outcome| !matches!(outcome.action, CopyAction::Orphaned { .. }));
1327 for entry in linking {
1328 let slot = landed
1329 .next()
1330 .filter(|outcome| outcome.source == entry.item)
1331 .and_then(|outcome| outcome.action.link_slot())
1332 .expect("every item a copy landed is reported, in the order it landed");
1333 *slot = Some(self.linked(destination, entry, journal).await?);
1334 }
1335 Ok(())
1336 }
1337
1338 /// Write one landed item's `onetaskgraph.copies` entry for the destination, when it
1339 /// does not already say this, and say what was done.
1340 ///
1341 /// An item found by its own origin records nothing: that correspondence is already
1342 /// recorded, on the destination item. A source that cannot hold the entry — it has no
1343 /// write side, or refuses the narrow metadata write — leaves the copy as it landed and
1344 /// is reported so; any other failure fails the copy, which then undoes every write it
1345 /// made, this one's journal entry included.
1346 async fn linked(
1347 &self,
1348 destination: &ResolvedSource,
1349 entry: Linking,
1350 journal: &mut Journal,
1351 ) -> Result<CopyLink, EngineError> {
1352 if entry.found == Some(Found::Origin) {
1353 return Ok(CopyLink::Unchanged);
1354 }
1355 let wanted = Value::String(entry.destination.to_string());
1356 // Only the entries that are links are kept. Anything else under the key — a value
1357 // that is not an object, or an entry a person spelled wrong — is not a link this or
1358 // any later copy can read, so it is replaced by what is, exactly as an entry naming
1359 // another item is rewritten.
1360 let mut links = well_formed(entry.held.as_ref());
1361 if links.get(destination.name().as_str()) == Some(&wanted) {
1362 return Ok(CopyLink::Unchanged);
1363 }
1364 let source = self.readable(&entry.item.source)?;
1365 if !source.source().writes().is_supported() {
1366 return Ok(CopyLink::Unrecorded);
1367 }
1368 links.insert(destination.name().to_string(), wanted);
1369 // Journalled before the write, for the reason a destination write is: a write that
1370 // stopped part way is only put back by what this journal holds.
1371 journal.links.push(Unlink {
1372 item: entry.item.clone(),
1373 level: entry.level,
1374 before: entry.held,
1375 });
1376 match set_link(
1377 source,
1378 entry.level,
1379 &entry.item.native,
1380 &Value::Object(links),
1381 )
1382 .await
1383 {
1384 Ok(true) => Ok(CopyLink::Recorded),
1385 // Nothing was written: the source no longer holds the item, or said it cannot
1386 // hold the key.
1387 Ok(false) | Err(SourceError::Refused { .. }) => {
1388 journal.links.pop();
1389 Ok(CopyLink::Unrecorded)
1390 }
1391 // §4.18 once let a plugin refuse any key in this namespace as malformed, and one
1392 // written then still may. Its answer says nothing about whether it wrote, so the
1393 // journal entry stays and a later failure puts the item back regardless.
1394 Err(SourceError::Malformed { .. }) => Ok(CopyLink::Unrecorded),
1395 Err(error) => Err(refused(source, error)),
1396 }
1397 }
1398
1399 /// Put one copied item's `onetaskgraph.copies` back to what it held before this copy.
1400 ///
1401 /// A key it held is written back through the same narrow write. A key it did not hold
1402 /// cannot be removed through that write, so the item is read as it now is and written
1403 /// back whole, without the key, through the source's own write interface — the same
1404 /// way an overwritten destination item is restored.
1405 async fn unlink(&self, entry: &Unlink) -> Result<(), SourceError> {
1406 let as_source_error = |error: EngineError| match error {
1407 EngineError::SourceRefused { error, .. } => error,
1408 other => SourceError::Refused {
1409 message: other.to_string(),
1410 },
1411 };
1412 let source = self.readable(&entry.item.source).map_err(as_source_error)?;
1413 if let Some(before) = &entry.before {
1414 return set_link(source, entry.level, &entry.item.native, before)
1415 .await
1416 .map(|_| ());
1417 }
1418 let Some(mut held) = self
1419 .prior(source, entry.level, &entry.item.native)
1420 .await
1421 .map_err(as_source_error)?
1422 else {
1423 return Ok(());
1424 };
1425 metadata_of(&mut held.item).remove(MetadataKey::COPIES_KEY);
1426 restore(source, &entry.item.native, &held).await
1427 }
1428
1429 /// The destination source, once it is established it exists and can be written.
1430 fn writable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
1431 let name = self.known(name)?;
1432 if let Some(unavailable) = self.unavailable().find(|source| source.name() == &name) {
1433 return Err(EngineError::DestinationUnavailable {
1434 name: name.to_string(),
1435 error: unavailable.error().clone(),
1436 });
1437 }
1438 let source = self
1439 .ready()
1440 .find(|source| source.name() == &name)
1441 .ok_or(EngineError::NoSources)?;
1442 if !source.source().writes().is_supported() {
1443 return Err(EngineError::NotWritable {
1444 name: name.to_string(),
1445 kind: source.kind().to_owned(),
1446 });
1447 }
1448 Ok(source)
1449 }
1450
1451 /// Copy each project and what travels with it: every task in it, the members a member
1452 /// copy names, or nothing at all.
1453 ///
1454 /// Every item of every project is read and its target decided **before any of them is
1455 /// written, and each exactly once.** That is what lets a member copy refuse an edge it
1456 /// cannot resolve while the destination is still as it was found, and what stops a
1457 /// whole copy resolving one target twice — once to learn the ids the project's own
1458 /// edges name, and again to land it — which against a hosted destination is a read per
1459 /// item spent on an answer the command already had.
1460 async fn copy_projects(
1461 &self,
1462 destination: &ResolvedSource,
1463 request: &CopyRequest,
1464 projects: &[(GlobalId, Vec<GlobalId>)],
1465 unrecorded: &[GlobalId],
1466 running: &mut Running,
1467 journal: &mut Journal,
1468 ) -> Result<Vec<CopyOutcome>, EngineError> {
1469 let mut plans = Vec::new();
1470 for (id, members) in projects {
1471 let project = self.plan(destination, request, Level::Project, id).await?;
1472 let mut tasks = Vec::new();
1473 for member in members {
1474 tasks.push(self.plan(destination, request, Level::Task, member).await?);
1475 }
1476 plans.push((project, tasks));
1477 }
1478 for (project, tasks) in &plans {
1479 for item in std::iter::once(project).chain(tasks) {
1480 unrecorded_far_end(item, destination, unrecorded)?;
1481 }
1482 }
1483 let mut outcomes = Vec::new();
1484 for (project, tasks) in plans {
1485 outcomes.extend(
1486 self.copy_project(destination, request, project, tasks, running, journal)
1487 .await?,
1488 );
1489 }
1490 Ok(outcomes)
1491 }
1492
1493 /// Land one planned project, then the tasks planned beside it.
1494 async fn copy_project(
1495 &self,
1496 destination: &ResolvedSource,
1497 request: &CopyRequest,
1498 project: Planned,
1499 tasks: Vec<Planned>,
1500 running: &mut Running,
1501 journal: &mut Journal,
1502 ) -> Result<Vec<CopyOutcome>, EngineError> {
1503 let (carries_tasks, walks_orphans) = match request.scope {
1504 CopyScope::Projects { tasks } => (tasks, tasks),
1505 CopyScope::Members(_) => (true, false),
1506 CopyScope::Tasks | CopyScope::Documents => (false, false),
1507 };
1508 let id = project.source.clone();
1509 let members: Vec<GlobalId> = tasks.iter().map(|task| task.source.clone()).collect();
1510 // The project lands with every edge it can already resolve — its members' targets
1511 // are known from their plans — so a project that has not changed is not written at
1512 // all. What it *reports* is read the way it always has been: against the ids known
1513 // before its own members' were, unless every edge resolves and nothing differs.
1514 // Reading it any other way would move the word a repeat copy reports for a project
1515 // whose only difference is an edge, which is not this change's to move.
1516 let mut before_members = running.counterparts.clone();
1517 if let Target::Update { id: target, .. } = &project.target {
1518 before_members.insert(id.to_string(), target.clone());
1519 }
1520 for item in std::iter::once(&project).chain(&tasks) {
1521 if let Target::Update { id: target, .. } = &item.target {
1522 running
1523 .counterparts
1524 .insert(item.source.to_string(), target.clone());
1525 }
1526 }
1527 let unchanged = match &project.target {
1528 Target::Update { id: target, .. } => {
1529 let settled = mapped_edges(
1530 &project.edges,
1531 &id.source,
1532 destination,
1533 &running.resolvable,
1534 &running.counterparts,
1535 );
1536 let first = mapped_edges(
1537 &project.edges,
1538 &id.source,
1539 destination,
1540 &running.resolvable,
1541 &before_members,
1542 );
1543 let held = project.held.as_ref();
1544 let unchanged_with = |edges: &[Option<DependencyEdge>]| {
1545 !changes(
1546 held,
1547 &project,
1548 target,
1549 &None,
1550 &resolved(edges),
1551 &[],
1552 destination.name(),
1553 )
1554 };
1555 Some(
1556 (!settled.iter().any(Option::is_none) && unchanged_with(&settled))
1557 || unchanged_with(&first),
1558 )
1559 }
1560 Target::Create => None,
1561 };
1562 // A project this copy creates holds nothing it did not file itself, so nothing filed
1563 // under it can be an orphan and the destination is not walked for one.
1564 let created = matches!(project.target, Target::Create);
1565 let mut outcomes = self
1566 .copy_items(destination, request, vec![project], None, running, journal)
1567 .await?;
1568 if let (Some(unchanged), Some(landed), Some(via)) = (
1569 unchanged,
1570 outcomes[0].destination().cloned(),
1571 outcomes[0].action.found_by(),
1572 ) {
1573 let link = outcomes[0].action.link();
1574 outcomes[0].action = if unchanged {
1575 CopyAction::Unchanged {
1576 destination: landed,
1577 via,
1578 link,
1579 }
1580 } else {
1581 CopyAction::Updated {
1582 destination: landed,
1583 via,
1584 link,
1585 }
1586 };
1587 }
1588 if !carries_tasks {
1589 return Ok(outcomes);
1590 }
1591 // `None` when a dry run would have created the project: nothing was written, so
1592 // there is no destination project id to file the tasks under. Every task is still
1593 // read and still reported, because that is what a dry run is for.
1594 let filed = outcomes.first().and_then(CopyOutcome::destination).cloned();
1595 outcomes.extend(
1596 self.copy_items(
1597 destination,
1598 request,
1599 tasks,
1600 filed.as_ref().map(|project| project.native.clone()),
1601 running,
1602 journal,
1603 )
1604 .await?,
1605 );
1606 // A member copy was told which members it carries, so it cannot tell a member it
1607 // left out from one the source no longer holds, and does not walk for either.
1608 if walks_orphans
1609 && !created
1610 && let Some(filed) = filed
1611 {
1612 outcomes.extend(
1613 self.orphans(destination, &id, &filed.native, &members)
1614 .await?,
1615 );
1616 }
1617 Ok(outcomes)
1618 }
1619
1620 /// The members a member copy names, per project and in the order they were named, and
1621 /// every member it does not name that records no destination id.
1622 ///
1623 /// Settled before anything is read at the destination. A named id that is a member of
1624 /// none of the projects is refused here, naming it. An unnamed member whose recorded
1625 /// origin names the destination joins the copied set with that id already known, so an
1626 /// edge to it resolves exactly as an edge to a copied item does, and costs no read. An
1627 /// unnamed member recording none is returned, so an edge to it is refused before
1628 /// anything is written — see [`unrecorded_far_end`].
1629 async fn named_members(
1630 &self,
1631 destination: &ResolvedSource,
1632 projects: &[GlobalId],
1633 named: &CopyItems,
1634 running: &mut Running,
1635 ) -> Result<(Vec<(GlobalId, Vec<GlobalId>)>, Vec<GlobalId>), EngineError> {
1636 let mut carried = Vec::new();
1637 let mut unrecorded = Vec::new();
1638 for project in projects {
1639 let held = self.project_member_tasks(project).await?;
1640 let mut members: Vec<GlobalId> = Vec::new();
1641 for id in named.as_slice() {
1642 if held.iter().any(|task| &task.id == id) && !members.contains(id) {
1643 members.push(id.clone());
1644 }
1645 }
1646 for task in held {
1647 if members.contains(&task.id) {
1648 continue;
1649 }
1650 match origin_of(&task.item.metadata)
1651 .filter(|origin| &origin.source == destination.name())
1652 {
1653 Some(origin) => {
1654 running
1655 .counterparts
1656 .insert(task.id.to_string(), origin.native);
1657 running.resolvable.push(task.id);
1658 }
1659 None => unrecorded.push(task.id),
1660 }
1661 }
1662 running.resolvable.push(project.clone());
1663 running.resolvable.extend(members.iter().cloned());
1664 carried.push((project.clone(), members));
1665 }
1666 if let Some(stray) = named
1667 .as_slice()
1668 .iter()
1669 .find(|id| !carried.iter().any(|(_, members)| members.contains(id)))
1670 {
1671 return Err(EngineError::NotAMember {
1672 id: stray.clone(),
1673 projects: projects.to_vec(),
1674 });
1675 }
1676 Ok((carried, unrecorded))
1677 }
1678
1679 /// Every task the source holds in `project`, by qualified id.
1680 async fn project_members(&self, project: &GlobalId) -> Result<Vec<GlobalId>, EngineError> {
1681 Ok(self
1682 .project_member_tasks(project)
1683 .await?
1684 .into_iter()
1685 .map(|task| task.id)
1686 .collect())
1687 }
1688
1689 /// Every task the source holds in `project`, as the source reported it.
1690 ///
1691 /// The ids alone are what a copy files under a project; the whole task is what a
1692 /// document's references need, because the location a reference names is a field of it.
1693 async fn project_member_tasks(
1694 &self,
1695 project: &GlobalId,
1696 ) -> Result<Vec<Qualified<Task>>, EngineError> {
1697 let mut request = TaskRequest {
1698 sources: vec![project.source.clone()],
1699 filters: Filters::default(),
1700 project: ProjectSelector::Qualified(project.clone()),
1701 priorities: Vec::new(),
1702 commented_since: None,
1703 metadata: Vec::new(),
1704 origin: None,
1705 paging: Paging {
1706 limit: PROJECT_PAGE,
1707 token: None,
1708 },
1709 };
1710 let mut members = Vec::new();
1711 // Pages by this engine's own token rather than by a source cursor, and the
1712 // asymmetry with the three walks below is deliberate. A source answering two
1713 // cursors with each other advances on every page, so `unrepeated` under the list
1714 // verb never fires; the cycle shows only as a token handed back unchanged, which
1715 // is what this loop pages by. Point it at the source and a project copy spins.
1716 //
1717 // No page bound here for the same reason: `Engine::tasks` merges under the budget
1718 // it was asked, so nothing longer than `PROJECT_PAGE` can arrive. `fits` in
1719 // `fetch::walk` refuses the page a source can really overrun, its own, while
1720 // these very members are read.
1721 let misbehaved = |error| EngineError::SourceRefused {
1722 name: project.source.to_string(),
1723 error,
1724 };
1725 loop {
1726 let asked = request.paging.token.clone();
1727 let response = self.tasks(&request).await?;
1728 if let Some(failure) = response.errors.first() {
1729 return Err(EngineError::SourceRefused {
1730 name: failure.source.to_string(),
1731 error: failure.error.clone(),
1732 });
1733 }
1734 unrepeated(
1735 response.next.as_ref(),
1736 asked.as_ref(),
1737 "the tasks of a project were being read for a copy",
1738 )
1739 .map_err(misbehaved)?;
1740 members.extend(response.items);
1741 match response.next {
1742 Some(token) => request.paging.token = Some(token),
1743 None => return Ok(members),
1744 }
1745 }
1746 }
1747
1748 /// Every document the source holds in `project`, as the source reported it.
1749 ///
1750 /// Paged by this engine's own token for the reason the member walk above is, and the
1751 /// note there says why.
1752 async fn project_documents(
1753 &self,
1754 project: &GlobalId,
1755 ) -> Result<Vec<Qualified<Document>>, EngineError> {
1756 let mut request = DocumentRequest {
1757 sources: vec![project.source.clone()],
1758 filters: DocumentFilters::default(),
1759 project: ProjectSelector::Qualified(project.clone()),
1760 paging: Paging {
1761 limit: PROJECT_PAGE,
1762 token: None,
1763 },
1764 };
1765 let mut held = Vec::new();
1766 let misbehaved = |error| EngineError::SourceRefused {
1767 name: project.source.to_string(),
1768 error,
1769 };
1770 loop {
1771 let asked = request.paging.token.clone();
1772 let response = self.documents(&request).await?;
1773 if let Some(failure) = response.errors.first() {
1774 return Err(EngineError::SourceRefused {
1775 name: failure.source.to_string(),
1776 error: failure.error.clone(),
1777 });
1778 }
1779 unrepeated(
1780 response.next.as_ref(),
1781 asked.as_ref(),
1782 "the documents of a project were being read for a copy",
1783 )
1784 .map_err(misbehaved)?;
1785 held.extend(response.items);
1786 match response.next {
1787 Some(token) => request.paging.token = Some(token),
1788 None => return Ok(held),
1789 }
1790 }
1791 }
1792
1793 /// Point every reference the documents of this copy hold at the destination's own
1794 /// records, and say how many it could not.
1795 ///
1796 /// A document copied out of a local Markdown store used to arrive naming absolute paths
1797 /// under one checkout on one machine, dead for the only reader the copy exists for,
1798 /// while the destination held its own record for every one of them the whole time. This
1799 /// is what closes that, and it is deliberately **not** a Markdown-link parser: the
1800 /// artifact that motivated it holds bare absolute paths inside backticks in a table
1801 /// cell, which `[text](target)` matching would have left exactly as it found them.
1802 ///
1803 /// The correspondence is re-established from what the *destination* records at
1804 /// [`GlobalId::ORIGIN_KEY`], not from the mapping this copy holds. That mapping is not
1805 /// available when it is needed: `project copy` and `document copy` are separate verbs,
1806 /// one invocation carries one [`CopyScope`], and a project copy carries no documents —
1807 /// so by the time the document is copied, the tasks were written by a process that has
1808 /// exited. Reading the destination is also what makes a document copied on its own
1809 /// work, which a same-run mapping never could.
1810 ///
1811 /// Tasks and projects are not touched. Only a document's content is rewritten, and
1812 /// nothing else about it changes.
1813 async fn rewrite_references(
1814 &self,
1815 destination: &ResolvedSource,
1816 planned: &mut [Planned],
1817 counts: &mut Counted,
1818 ) -> Result<(), EngineError> {
1819 // Read at the source, once per project rather than once per document: several
1820 // documents of one project is the ordinary case.
1821 let mut by_project: BTreeMap<String, Vec<Referent>> = BTreeMap::new();
1822 let mut named: Vec<Vec<Referent>> = Vec::new();
1823 for item in planned.iter() {
1824 named.push(self.named_referents(item, &mut by_project).await?);
1825 }
1826 // The destination is walked only for a copy that really names something, and only
1827 // for the interfaces those referents are read from.
1828 let mut levels: Vec<Level> = Vec::new();
1829 for referent in named.iter().flatten() {
1830 if !levels.contains(&referent.level) {
1831 levels.push(referent.level);
1832 }
1833 }
1834 if levels.is_empty() {
1835 return Ok(());
1836 }
1837 let counterparts = self.counterparts(destination, &levels).await?;
1838 for (item, referents) in planned.iter_mut().zip(named) {
1839 let Item::Document(document) = &mut item.item else {
1840 continue;
1841 };
1842 let Some(content) = &document.content else {
1843 continue;
1844 };
1845 let (rewritten, made) = substitute(content, &table_for(&referents, &counterparts));
1846 document.content = Some(rewritten);
1847 counts.add(made);
1848 }
1849 Ok(())
1850 }
1851
1852 /// The records of one document's own project whose location string its content really
1853 /// holds, as a whole reference.
1854 ///
1855 /// A document with no content, or with no project at the source, names nothing: the
1856 /// referent set is the document's own project, its tasks and its other documents, and
1857 /// there is no such set without a project.
1858 async fn named_referents(
1859 &self,
1860 item: &Planned,
1861 by_project: &mut BTreeMap<String, Vec<Referent>>,
1862 ) -> Result<Vec<Referent>, EngineError> {
1863 let Item::Document(document) = &item.item else {
1864 return Ok(Vec::new());
1865 };
1866 let (Some(content), Some(project)) =
1867 (document.content.as_deref(), document.project.as_ref())
1868 else {
1869 return Ok(Vec::new());
1870 };
1871 if content.is_empty() {
1872 return Ok(Vec::new());
1873 }
1874 let project = GlobalId::new(item.source.source.clone(), project.clone());
1875 let key = project.to_string();
1876 if !by_project.contains_key(&key) {
1877 let read = self.referents(&project).await?;
1878 by_project.insert(key.clone(), read);
1879 }
1880 Ok(by_project[&key]
1881 .iter()
1882 // A document does not name itself: the referent set is every *other* record
1883 // filed under the project. Told by the interface as well as the id, because an
1884 // id alone does not identify a record — a folder of Markdown filing `A.md`
1885 // under both `tasks/` and `documents/` is the ordinary case rather than the
1886 // contrived one, and excluding by id alone would drop that task from the set.
1887 .filter(|referent| referent.level != Level::Document || referent.id != item.source)
1888 .filter(|referent| holds(content, &referent.location))
1889 .cloned()
1890 .collect())
1891 }
1892
1893 /// Every record filed under one project at its own source, with the location string
1894 /// that source reports for it.
1895 ///
1896 /// The project record itself, every task filed under it, and every document filed under
1897 /// it. All three are reads this engine already knows how to make.
1898 async fn referents(&self, project: &GlobalId) -> Result<Vec<Referent>, EngineError> {
1899 let source = self.readable(&project.source)?;
1900 let mut referents = Vec::new();
1901 if let Some(held) = source
1902 .source()
1903 .get_project(&project.native)
1904 .await
1905 .map_err(|error| refused(source, error))?
1906 {
1907 note(
1908 &mut referents,
1909 project.clone(),
1910 Level::Project,
1911 held.location.as_ref(),
1912 &held.metadata,
1913 );
1914 }
1915 for task in self.project_member_tasks(project).await? {
1916 note(
1917 &mut referents,
1918 task.id,
1919 Level::Task,
1920 task.item.location.as_ref(),
1921 &task.item.metadata,
1922 );
1923 }
1924 for document in self.project_documents(project).await? {
1925 note(
1926 &mut referents,
1927 document.id,
1928 Level::Document,
1929 document.item.location.as_ref(),
1930 &document.item.metadata,
1931 );
1932 }
1933 Ok(referents)
1934 }
1935
1936 /// Walk the destination once for every record it holds at the levels named, one page
1937 /// at a time.
1938 ///
1939 /// One page is *read* at a time, and what is kept from each is three fields of each
1940 /// record — its id, the origin it records and the location the destination reports —
1941 /// never the page. That is more than [`Engine::scan`] keeps, and deliberately: the whole
1942 /// point of this walk is that one pass answers every document and every referent of the
1943 /// invocation, so what it learns has to outlive the page it learned it from. Nothing is
1944 /// written down and the index is dropped with the call.
1945 ///
1946 /// The other difference from `scan` is the *answer*: this one keeps every match, so a
1947 /// destination holding two records for one work item is reported as ambiguous rather
1948 /// than resolved to the first.
1949 async fn counterparts(
1950 &self,
1951 destination: &ResolvedSource,
1952 levels: &[Level],
1953 ) -> Result<Counterparts, EngineError> {
1954 let mut found = Counterparts::default();
1955 for level in levels {
1956 // Every cursor this level has already been sent. `unrepeated` below catches a
1957 // source that hands back the cursor it was just given, and its own note says
1958 // why it catches no more than that: a source cycling through two cursors
1959 // advances on every page, and seeing it needs memory the walks that share that
1960 // helper do not keep. This walk does keep memory — it is building an index that
1961 // outlives each page — so here the memory exists and the cycle is caught. The
1962 // walks one level up catch the same defect as a page token handed back
1963 // unchanged; this one pages by the source's own cursor and has no level above
1964 // it, so nothing else would.
1965 let mut asked_before: BTreeSet<String> = BTreeSet::new();
1966 let mut cursor: Option<Cursor> = None;
1967 loop {
1968 if let Some(next) = &cursor
1969 && !asked_before.insert(next.0.clone())
1970 {
1971 return Err(refused(
1972 destination,
1973 SourceError::Malformed {
1974 message: "the source returned a cursor it had already been \
1975 given while the destination was being walked for the \
1976 records a document's references name, so the walk \
1977 would never end"
1978 .to_owned(),
1979 },
1980 ));
1981 }
1982 let asked = cursor.clone();
1983 let request = request_for(destination, cursor);
1984 let next = match level {
1985 Level::Task => {
1986 let page = destination
1987 .source()
1988 .query_tasks(&TaskQuery::default(), &request)
1989 .await
1990 .map_err(|error| refused(destination, error))?;
1991 fits(page.items.len(), request.limit)
1992 .map_err(|error| refused(destination, error))?;
1993 for task in &page.items {
1994 found.note(*level, &task.id, task.location.as_ref(), &task.metadata);
1995 }
1996 page.next
1997 }
1998 Level::Project => {
1999 let page = destination
2000 .source()
2001 .query_projects(&ProjectQuery::default(), &request)
2002 .await
2003 .map_err(|error| refused(destination, error))?;
2004 fits(page.items.len(), request.limit)
2005 .map_err(|error| refused(destination, error))?;
2006 for project in &page.items {
2007 found.note(
2008 *level,
2009 &project.id,
2010 project.location.as_ref(),
2011 &project.metadata,
2012 );
2013 }
2014 page.next
2015 }
2016 Level::Document => {
2017 let page = destination
2018 .source()
2019 .query_documents(&DocumentQuery::default(), &request)
2020 .await
2021 .map_err(|error| refused(destination, error))?;
2022 fits(page.items.len(), request.limit)
2023 .map_err(|error| refused(destination, error))?;
2024 for document in &page.items {
2025 found.note(
2026 *level,
2027 &document.id,
2028 document.location.as_ref(),
2029 &document.metadata,
2030 );
2031 }
2032 page.next
2033 }
2034 };
2035 unrepeated(
2036 next.as_ref(),
2037 asked.as_ref(),
2038 "the destination was being walked for the records a document's \
2039 references name",
2040 )
2041 .map_err(|error| refused(destination, error))?;
2042 match next {
2043 Some(next) => cursor = Some(next),
2044 None => break,
2045 }
2046 }
2047 }
2048 Ok(found)
2049 }
2050
2051 /// Destination tasks filed under the copied project whose origin the source no longer
2052 /// holds.
2053 ///
2054 /// A copy never deletes, so each is left exactly as it is and reported.
2055 async fn orphans(
2056 &self,
2057 destination: &ResolvedSource,
2058 project: &GlobalId,
2059 at_destination: &NativeId,
2060 copied: &[GlobalId],
2061 ) -> Result<Vec<CopyOutcome>, EngineError> {
2062 let mut orphans = Vec::new();
2063 let mut cursor: Option<Cursor> = None;
2064 loop {
2065 let asked = cursor.clone();
2066 let request = request_for(destination, cursor);
2067 let page: Page<Task> = destination
2068 .source()
2069 .query_tasks(&TaskQuery::default(), &request)
2070 .await
2071 .map_err(|error| refused(destination, error))?;
2072 fits(page.items.len(), request.limit).map_err(|error| refused(destination, error))?;
2073 for task in &page.items {
2074 if task.project.as_ref() != Some(at_destination) {
2075 continue;
2076 }
2077 let Some(origin) = origin_of(&task.metadata) else {
2078 continue;
2079 };
2080 if origin.source != project.source || copied.contains(&origin) {
2081 continue;
2082 }
2083 orphans.push(CopyOutcome {
2084 source: origin,
2085 action: CopyAction::Orphaned {
2086 destination: GlobalId::new(destination.name().clone(), task.id.clone()),
2087 },
2088 });
2089 }
2090 unrepeated(
2091 page.next.as_ref(),
2092 asked.as_ref(),
2093 "the destination was being read for items the copy left behind",
2094 )
2095 .map_err(|error| refused(destination, error))?;
2096 match page.next {
2097 Some(next) => cursor = Some(next),
2098 None => return Ok(orphans),
2099 }
2100 }
2101 }
2102
2103 /// Resolve and write every item planned, holding back the ones whose edges are not
2104 /// resolvable yet.
2105 ///
2106 /// An edge between two items of one copy can point at a member whose destination id
2107 /// does not exist until it has been created, so the item that points at it lands
2108 /// without that edge and is handed to `deferred`. [`Engine::repair`] finishes it once
2109 /// the *whole* request has landed — not once this call has, because the far end may
2110 /// be in another project of the same command.
2111 async fn copy_items(
2112 &self,
2113 destination: &ResolvedSource,
2114 request: &CopyRequest,
2115 mut planned: Vec<Planned>,
2116 project: Option<NativeId>,
2117 running: &mut Running,
2118 journal: &mut Journal,
2119 ) -> Result<Vec<CopyOutcome>, EngineError> {
2120 // Only a document's own content names other records, and only once every document
2121 // of this call has been read: the destination is walked once for all of them, and
2122 // the content the rest of this call lands is the rewritten one — which is what
2123 // makes a repeat copy of an already-rewritten document report `unchanged`.
2124 if planned
2125 .iter()
2126 .any(|item| item.item.level() == Level::Document)
2127 {
2128 self.rewrite_references(destination, &mut planned, &mut running.references)
2129 .await?;
2130 }
2131
2132 for item in &planned {
2133 if let Target::Update { id, .. } = &item.target {
2134 running
2135 .counterparts
2136 .insert(item.source.to_string(), id.clone());
2137 }
2138 if let Item::Task(task) = &item.item {
2139 running.delivers_rewritten +=
2140 members_named(&task.delivers, &item.source.source, &running.resolvable);
2141 }
2142 }
2143
2144 // Resolved once per item, and used by both passes: the repair pass writes the
2145 // same item again, and re-deriving this there could file it somewhere else.
2146 let mut filed = Vec::new();
2147 for item in &planned {
2148 filed.push(
2149 self.filed(destination, item, project.clone(), running)
2150 .await?,
2151 );
2152 }
2153
2154 let mut outcomes = Vec::new();
2155 let mut unresolved = Vec::new();
2156 let mut priors = Vec::new();
2157 for (index, item) in planned.iter().enumerate() {
2158 let edges = mapped_edges(
2159 &item.edges,
2160 &item.source.source,
2161 destination,
2162 &running.resolvable,
2163 &running.counterparts,
2164 );
2165 let delivers = delivers_of(
2166 item,
2167 destination,
2168 &running.resolvable,
2169 &running.counterparts,
2170 );
2171 if edges.iter().any(Option::is_none) || delivers.iter().any(Option::is_none) {
2172 unresolved.push(index);
2173 }
2174 let resolvable_delivers = resolved_entries(&delivers);
2175 let (outcome, prior) = self
2176 .land(
2177 destination,
2178 request,
2179 item,
2180 filed[index].clone(),
2181 Pointing {
2182 edges: &edges,
2183 delivers: &resolvable_delivers,
2184 },
2185 journal,
2186 )
2187 .await?;
2188 if let Some(id) = outcome.destination() {
2189 running
2190 .counterparts
2191 .insert(item.source.to_string(), id.native.clone());
2192 if !request.dry_run {
2193 running.linking.push(Linking {
2194 item: item.source.clone(),
2195 level: item.item.level(),
2196 found: item.target.found(),
2197 destination: id.clone(),
2198 held: described(&item.item)
2199 .1
2200 .get(MetadataKey::COPIES_KEY)
2201 .cloned(),
2202 });
2203 }
2204 }
2205 if !request.dry_run
2206 && let (Item::Task(task), Some(landed)) = (&item.item, outcome.destination())
2207 {
2208 running.landed.push(LandedTask {
2209 destination: landed.clone(),
2210 origin: item.source.source.clone(),
2211 delivers: task.delivers.clone(),
2212 before: match prior.as_ref().map(|prior| &prior.item) {
2213 Some(Item::Task(held)) => targets(&held.delivers, destination.name()),
2214 _ => Vec::new(),
2215 },
2216 category: task.status.category,
2217 });
2218 }
2219 outcomes.push(outcome);
2220 priors.push(prior);
2221 }
2222
2223 if !request.dry_run {
2224 for (index, item) in planned.into_iter().enumerate() {
2225 if !unresolved.contains(&index) {
2226 continue;
2227 }
2228 // Every item a copy that is not a dry run lands has a destination id: the
2229 // one outcome without one is a dry run that would have created, and this
2230 // block does not run for a dry run.
2231 let id = outcomes[index]
2232 .destination()
2233 .expect("a copy that writes lands every item it planned")
2234 .clone();
2235 running.deferred.push(Deferred {
2236 item,
2237 filed: filed[index].clone(),
2238 destination: id.native,
2239 prior: priors[index].clone(),
2240 });
2241 }
2242 }
2243 Ok(outcomes)
2244 }
2245
2246 /// Read one item and its forward edges, and decide where it is going.
2247 async fn plan(
2248 &self,
2249 destination: &ResolvedSource,
2250 request: &CopyRequest,
2251 kind: Level,
2252 id: &GlobalId,
2253 ) -> Result<Planned, EngineError> {
2254 let source = self.readable(&id.source)?;
2255 if kind == Level::Document {
2256 documentary(source)?;
2257 }
2258 let item = match kind {
2259 Level::Task => source
2260 .source()
2261 .get_task(&id.native)
2262 .await
2263 .map_err(|error| refused(source, error))?
2264 .map(|task| Item::Task(Box::new(task))),
2265 Level::Project => source
2266 .source()
2267 .get_project(&id.native)
2268 .await
2269 .map_err(|error| refused(source, error))?
2270 .map(|project| Item::Project(Box::new(project))),
2271 Level::Document => source
2272 .source()
2273 .get_document(&id.native)
2274 .await
2275 .map_err(|error| refused(source, error))?
2276 .map(|document| Item::Document(Box::new(document))),
2277 }
2278 .ok_or_else(|| EngineError::NoSuchItem { id: id.to_string() })?;
2279 // Before the destination is read, and so before it is written: a destination that
2280 // holds no priority has nowhere to put one, and every item of a copy is planned before
2281 // any of them lands. A task carrying `none` passes and writes exactly as it always did.
2282 if let Item::Task(task) = &item {
2283 holds_priority(destination, &id.to_string(), task.priority)?;
2284 }
2285 let edges = forward_edges(source, &id.native, item.level()).await?;
2286 let (target, held) = self.target(destination, request, id, &item).await?;
2287 Ok(Planned {
2288 source: id.clone(),
2289 item,
2290 edges,
2291 target,
2292 held,
2293 })
2294 }
2295
2296 /// Which destination item this one corresponds to, by the link the item records, the
2297 /// two origin rules and the caller's escape, and what the destination holds there.
2298 async fn target(
2299 &self,
2300 destination: &ResolvedSource,
2301 request: &CopyRequest,
2302 id: &GlobalId,
2303 item: &Item,
2304 ) -> Result<(Target, Option<Prior>), EngineError> {
2305 let (title, metadata) = described(item);
2306 // The link first, and it is believed only when the item it names still names this
2307 // one back: a person who re-pointed that item's origin has said it is not this
2308 // one's counterpart any more, so the link is ignored and rewritten to whatever the
2309 // rules below find.
2310 if let Some(link) = link_of(metadata, destination.name()) {
2311 match self.prior(destination, item.level(), &link.native).await? {
2312 Some(held) if origin_of(described(&held.item).1).as_ref() == Some(id) => {
2313 return Ok((
2314 Target::Update {
2315 id: link.native,
2316 found: Found::Link,
2317 },
2318 Some(held),
2319 ));
2320 }
2321 Some(_) => {}
2322 None if request.recreate => {}
2323 None => {
2324 return Err(EngineError::StaleLink {
2325 item: id.to_string(),
2326 link: link.to_string(),
2327 });
2328 }
2329 }
2330 }
2331 if let Some(origin) = origin_of(metadata)
2332 && &origin.source == destination.name()
2333 {
2334 // The read that says the origin still names something is the read of what it
2335 // holds, so rule 1 costs one round trip rather than two.
2336 if let Some(held) = self
2337 .prior(destination, item.level(), &origin.native)
2338 .await?
2339 {
2340 return Ok((
2341 Target::Update {
2342 id: origin.native,
2343 found: Found::Origin,
2344 },
2345 Some(held),
2346 ));
2347 }
2348 if !request.recreate {
2349 return Err(EngineError::StaleOrigin {
2350 item: id.to_string(),
2351 origin: origin.to_string(),
2352 });
2353 }
2354 }
2355 if let Some(found) = self
2356 .scan(destination, item.level(), &Wanted::Origin(id.to_string()))
2357 .await?
2358 {
2359 let held = self.prior(destination, item.level(), &found).await?;
2360 return Ok((
2361 Target::Update {
2362 id: found,
2363 found: Found::Scan,
2364 },
2365 held,
2366 ));
2367 }
2368 let wanted = match &request.match_by {
2369 Some(MatchBy::Title) => Some(Wanted::Title(title.to_owned())),
2370 Some(MatchBy::Metadata(key)) => metadata
2371 .get(key)
2372 .map(|value| Wanted::Metadata(key.clone(), value.clone())),
2373 None => None,
2374 };
2375 if let Some(wanted) = wanted
2376 && let Some(found) = self.scan(destination, item.level(), &wanted).await?
2377 {
2378 let held = self.prior(destination, item.level(), &found).await?;
2379 return Ok((
2380 Target::Update {
2381 id: found,
2382 found: Found::Match,
2383 },
2384 held,
2385 ));
2386 }
2387 Ok((Target::Create, None))
2388 }
2389
2390 /// Walk the destination one page at a time, looking for `wanted`.
2391 ///
2392 /// One page is held at a time and nothing is written down, which is the same bound
2393 /// every other compensation in this engine works under.
2394 async fn scan(
2395 &self,
2396 destination: &ResolvedSource,
2397 kind: Level,
2398 wanted: &Wanted,
2399 ) -> Result<Option<NativeId>, EngineError> {
2400 let mut cursor: Option<Cursor> = None;
2401 let tasks = wanted.task_query(&destination.source().capabilities());
2402 loop {
2403 let asked = cursor.clone();
2404 let request = request_for(destination, cursor);
2405 let next = match kind {
2406 Level::Task => {
2407 let page = destination
2408 .source()
2409 .query_tasks(&tasks, &request)
2410 .await
2411 .map_err(|error| refused(destination, error))?;
2412 fits(page.items.len(), request.limit)
2413 .map_err(|error| refused(destination, error))?;
2414 for task in &page.items {
2415 if wanted.found(&task.title, &task.metadata) {
2416 return Ok(Some(task.id.clone()));
2417 }
2418 }
2419 page.next
2420 }
2421 Level::Project => {
2422 let page = destination
2423 .source()
2424 .query_projects(&ProjectQuery::default(), &request)
2425 .await
2426 .map_err(|error| refused(destination, error))?;
2427 fits(page.items.len(), request.limit)
2428 .map_err(|error| refused(destination, error))?;
2429 for project in &page.items {
2430 if wanted.found(&project.title, &project.metadata) {
2431 return Ok(Some(project.id.clone()));
2432 }
2433 }
2434 page.next
2435 }
2436 Level::Document => {
2437 let page = destination
2438 .source()
2439 .query_documents(&DocumentQuery::default(), &request)
2440 .await
2441 .map_err(|error| refused(destination, error))?;
2442 fits(page.items.len(), request.limit)
2443 .map_err(|error| refused(destination, error))?;
2444 for document in &page.items {
2445 if wanted.found(&document.title, &document.metadata) {
2446 return Ok(Some(document.id.clone()));
2447 }
2448 }
2449 page.next
2450 }
2451 };
2452 unrepeated(
2453 next.as_ref(),
2454 asked.as_ref(),
2455 "the destination was being scanned for the item to update",
2456 )
2457 .map_err(|error| refused(destination, error))?;
2458 match next {
2459 Some(next) => cursor = Some(next),
2460 None => return Ok(None),
2461 }
2462 }
2463 }
2464
2465 /// Write one planned item, or say what a dry run would have done.
2466 ///
2467 /// Answers with what the destination held there beforehand as well, which is what
2468 /// makes an item written twice restorable to what it was rather than to what this
2469 /// copy's first pass left.
2470 async fn land(
2471 &self,
2472 destination: &ResolvedSource,
2473 request: &CopyRequest,
2474 item: &Planned,
2475 project: Option<NativeId>,
2476 pointing: Pointing<'_>,
2477 journal: &mut Journal,
2478 ) -> Result<(CopyOutcome, Option<Prior>), EngineError> {
2479 let Pointing { edges, delivers } = pointing;
2480 let (target, found) = match &item.target {
2481 Target::Update { id, found } => (Some(id.clone()), Some(*found)),
2482 Target::Create => (None, None),
2483 };
2484 // Every link is settled once the whole copy has landed, by `Engine::link`, so none
2485 // is known here.
2486 let reached = |destination: Option<GlobalId>, changed: bool| match (found, destination) {
2487 (Some(via), Some(destination)) if changed => CopyAction::Updated {
2488 destination,
2489 via: via.via(),
2490 link: None,
2491 },
2492 (Some(via), Some(destination)) => CopyAction::Unchanged {
2493 destination,
2494 via: via.via(),
2495 link: None,
2496 },
2497 (_, destination) => CopyAction::Created {
2498 destination,
2499 via: NoCounterpart::Created,
2500 link: None,
2501 },
2502 };
2503 // The one read of the destination item, made where its target was found, used to
2504 // decide whether the write changes anything and — if the copy cannot finish — to
2505 // put that item back.
2506 let prior = item.held.clone();
2507 let edges = resolved(edges);
2508 let qualified = |native: NativeId| GlobalId::new(destination.name().clone(), native);
2509 if let Some(id) = &target
2510 && !changes(
2511 prior.as_ref(),
2512 item,
2513 id,
2514 &project,
2515 &edges,
2516 delivers,
2517 destination.name(),
2518 )
2519 {
2520 return Ok((
2521 CopyOutcome {
2522 source: item.source.clone(),
2523 action: reached(Some(qualified(id.clone())), false),
2524 },
2525 prior,
2526 ));
2527 }
2528 if request.dry_run {
2529 return Ok((
2530 CopyOutcome {
2531 source: item.source.clone(),
2532 // Null only for a create here: nothing was created, so there is no id to
2533 // report.
2534 action: reached(target.map(qualified), true),
2535 },
2536 prior,
2537 ));
2538 }
2539 let written = qualified(
2540 self.write(
2541 destination,
2542 item,
2543 target,
2544 project,
2545 &edges,
2546 delivers,
2547 prior.clone(),
2548 journal,
2549 )
2550 .await?,
2551 );
2552 Ok((
2553 CopyOutcome {
2554 source: item.source.clone(),
2555 action: reached(Some(written), true),
2556 },
2557 prior,
2558 ))
2559 }
2560
2561 /// Which destination project this item is filed under, when it is filed at all.
2562 ///
2563 /// A task copied as part of a project copy is filed under that project's counterpart,
2564 /// which the copy has just established. A task copied on its own has to find it.
2565 async fn filed(
2566 &self,
2567 destination: &ResolvedSource,
2568 item: &Planned,
2569 project: Option<NativeId>,
2570 running: &mut Running,
2571 ) -> Result<Option<NativeId>, EngineError> {
2572 match (&item.item, project) {
2573 (Item::Task(task), None) => {
2574 self.counterpart(destination, item, task.project.as_ref(), running)
2575 .await
2576 }
2577 (Item::Document(document), None) => {
2578 self.counterpart(destination, item, document.project.as_ref(), running)
2579 .await
2580 }
2581 (Item::Task(_) | Item::Document(_), filed) => Ok(filed),
2582 (Item::Project(_), _) => Ok(None),
2583 }
2584 }
2585
2586 /// The destination project this task's own project corresponds to, when there is one.
2587 ///
2588 /// A task copied on its own keeps its source's project id when the destination holds
2589 /// no counterpart: the field is opaque to this engine, and dropping it would lose
2590 /// what the source said.
2591 ///
2592 /// Looked for once per source project per command. A project this command itself
2593 /// copied is already known, and one an earlier task of this command was filed under
2594 /// was already looked for; walking the destination again for either would spend a
2595 /// scan per task on an answer the command holds. A project whose own link names its
2596 /// counterpart there is found by that link — see [`Engine::linked_project`] — and the
2597 /// destination is walked only when it does not.
2598 async fn counterpart(
2599 &self,
2600 destination: &ResolvedSource,
2601 item: &Planned,
2602 project: Option<&NativeId>,
2603 running: &mut Running,
2604 ) -> Result<Option<NativeId>, EngineError> {
2605 let Some(project) = project else {
2606 return Ok(None);
2607 };
2608 let qualified = GlobalId::new(item.source.source.clone(), project.clone()).to_string();
2609 if let Some(landed) = running.counterparts.get(&qualified) {
2610 return Ok(Some(landed.clone()));
2611 }
2612 if let Some(looked) = running.filings.get(&qualified) {
2613 return Ok(looked.clone());
2614 }
2615 let found = match self
2616 .linked_project(destination, &item.source.source, project, &qualified)
2617 .await?
2618 {
2619 Some(linked) => Some(linked),
2620 None => {
2621 self.scan(
2622 destination,
2623 Level::Project,
2624 &Wanted::Origin(qualified.clone()),
2625 )
2626 .await?
2627 }
2628 };
2629 let filed = Some(found.unwrap_or_else(|| project.clone()));
2630 running.filings.insert(qualified, filed.clone());
2631 Ok(filed)
2632 }
2633
2634 /// The destination project a source project's own `onetaskgraph.copies` entry names,
2635 /// when the project it names there still records that source project as its origin.
2636 ///
2637 /// The filing lookup's half of following a link: one read of the project at its source
2638 /// and one of its counterpart by id, in place of a walk of the destination's projects.
2639 /// A link that names nothing there, or names a project somebody re-pointed, is not a
2640 /// refusal here as it is for a copy's own target — this only decides where an item is
2641 /// filed — so the walk answers instead, exactly as it did before there were links.
2642 async fn linked_project(
2643 &self,
2644 destination: &ResolvedSource,
2645 source: &SourceName,
2646 project: &NativeId,
2647 qualified: &str,
2648 ) -> Result<Option<NativeId>, EngineError> {
2649 let source = self.readable(source)?;
2650 if !source.source().capabilities().projects.is_native() {
2651 return Ok(None);
2652 }
2653 let Some(held) = source
2654 .source()
2655 .get_project(project)
2656 .await
2657 .map_err(|error| refused(source, error))?
2658 else {
2659 return Ok(None);
2660 };
2661 let Some(link) = link_of(&held.metadata, destination.name()) else {
2662 return Ok(None);
2663 };
2664 let there = destination
2665 .source()
2666 .get_project(&link.native)
2667 .await
2668 .map_err(|error| refused(destination, error))?;
2669 Ok(there
2670 .filter(|there| {
2671 origin_of(&there.metadata).is_some_and(|origin| origin.to_string() == qualified)
2672 })
2673 .map(|_| link.native))
2674 }
2675
2676 /// What the destination holds at one id, item and forward edges together.
2677 ///
2678 /// One read for both purposes it serves — deciding whether a write changes anything,
2679 /// and putting the item back if the copy cannot finish — because a second read of the
2680 /// same item is a second round trip against a hosted destination for nothing.
2681 async fn prior(
2682 &self,
2683 destination: &ResolvedSource,
2684 kind: Level,
2685 id: &NativeId,
2686 ) -> Result<Option<Prior>, EngineError> {
2687 let held = match kind {
2688 Level::Task => destination
2689 .source()
2690 .get_task(id)
2691 .await
2692 .map_err(|error| refused(destination, error))?
2693 .map(|task| Item::Task(Box::new(task))),
2694 Level::Project => destination
2695 .source()
2696 .get_project(id)
2697 .await
2698 .map_err(|error| refused(destination, error))?
2699 .map(|project| Item::Project(Box::new(project))),
2700 Level::Document => destination
2701 .source()
2702 .get_document(id)
2703 .await
2704 .map_err(|error| refused(destination, error))?
2705 .map(|document| Item::Document(Box::new(document))),
2706 };
2707 let Some(item) = held else {
2708 return Ok(None);
2709 };
2710 let edges = forward_edges(destination, id, kind).await?;
2711 Ok(Some(Prior { item, edges }))
2712 }
2713
2714 /// Hand one item to the destination's own write interface, recording how to take it
2715 /// back.
2716 // llmlint: ignore[suppressions_justified] A write is the item, where it is going, what
2717 // it is filed under, its edges, what was there before and the journal that records how
2718 // to put it back. Each is a distinct decision made by a different part of the copy, and
2719 // grouping them would only move the argument list to a constructor.
2720 #[allow(clippy::too_many_arguments)]
2721 async fn write(
2722 &self,
2723 destination: &ResolvedSource,
2724 item: &Planned,
2725 target: Option<NativeId>,
2726 project: Option<NativeId>,
2727 edges: &[DependencyEdge],
2728 delivers: &[TaskRef],
2729 prior: Option<Prior>,
2730 journal: &mut Journal,
2731 ) -> Result<NativeId, EngineError> {
2732 let created_kind = item.item.level();
2733 let suggested = target
2734 .clone()
2735 .unwrap_or_else(|| created_id(&item.item, project.as_ref()));
2736 // Settled before the journal takes `prior`, and from that same read: what the
2737 // destination holds at the origin key is what a copy-back leaves there, and what it
2738 // holds as `delivered_by` is what the item keeps.
2739 let origin = recorded(item, prior.as_ref());
2740 let landing = outgoing(item, suggested, project, &origin, delivers, prior.as_ref());
2741 // Recorded *before* the write rather than after it. A destination's own write is
2742 // several calls — `docs/plugin-protocol.md` §4.9 — and one of them failing leaves
2743 // the ones before it applied. No source can put those back, because only this
2744 // journal holds what was there; recorded after a successful write, an update that
2745 // stopped part way was the one way a copy could end and leave the destination
2746 // altered. A restore of an item the write never reached rewrites what is already
2747 // there, which costs one mutation and is what "either complete or it never
2748 // happened" is worth.
2749 if let (Some(id), Some(prior)) = (target.clone(), prior) {
2750 journal.record(Undo::Updated { id, prior });
2751 }
2752 let landed = match landing {
2753 Item::Task(task) => destination
2754 .source()
2755 .write_task(&ItemWrite {
2756 target: target.clone(),
2757 item: *task,
2758 depends_on: edges.to_vec(),
2759 })
2760 .await
2761 .map_err(|error| refused(destination, error))?,
2762 Item::Project(project) => destination
2763 .source()
2764 .write_project(&ItemWrite {
2765 target: target.clone(),
2766 item: *project,
2767 depends_on: edges.to_vec(),
2768 })
2769 .await
2770 .map_err(|error| refused(destination, error))?,
2771 // No edges, and that is the contract: a document takes part in no dependency
2772 // graph, so there is nothing here for `depends_on` to carry.
2773 Item::Document(document) => destination
2774 .source()
2775 .write_document(&ItemWrite {
2776 target: target.clone(),
2777 item: *document,
2778 depends_on: Vec::new(),
2779 })
2780 .await
2781 .map_err(|error| refused(destination, error))?,
2782 };
2783 // A created item can only be journalled here: its id is what the write answers
2784 // with. A create that fails leaves nothing behind — §4.9 makes taking the item
2785 // back the source's own duty, because a write that refused must not leave an item
2786 // nobody asked for.
2787 if target.is_none() {
2788 journal.record(Undo::Created {
2789 kind: created_kind,
2790 id: landed.clone(),
2791 });
2792 }
2793 Ok(landed)
2794 }
2795
2796 /// A configured source that built, for reading an item out of.
2797 fn readable(&self, name: &SourceName) -> Result<&ResolvedSource, EngineError> {
2798 let name = self.known(name)?;
2799 if let Some(unavailable) = self.unavailable().find(|source| source.name() == &name) {
2800 return Err(EngineError::SourceRefused {
2801 name: name.to_string(),
2802 error: unavailable.error().clone(),
2803 });
2804 }
2805 self.ready()
2806 .find(|source| source.name() == &name)
2807 .ok_or(EngineError::NoSources)
2808 }
2809}
2810
2811/// How many tasks of a project are read at once while walking it.
2812const PROJECT_PAGE: std::num::NonZeroU32 = std::num::NonZeroU32::new(50).expect("50 is not zero");
2813
2814/// Each source's own running totals, in the order the sources were given.
2815///
2816/// A reading a source could not take is read as that source not metering: what a command
2817/// spent is a report about the work, and a failed reading must not become a failure of the
2818/// work itself. That holds when only one of a source's two readings failed as well, since
2819/// a difference needs both ends.
2820pub(super) async fn readings(sources: &[&ResolvedSource]) -> Vec<Option<Metering>> {
2821 let mut read = Vec::with_capacity(sources.len());
2822 for source in sources {
2823 read.push(source.source().metering().await.ok().flatten());
2824 }
2825 read
2826}
2827
2828/// What the sources spent between two readings of them, or `None` when none of them meters.
2829///
2830/// A source that does not meter contributes nothing and does not make the total zero: a
2831/// command whose sources all declined reports that it cannot say, not that it spent nothing.
2832/// A budget is keyed by its name and unit together, so two sources naming one budget add up
2833/// and two units of one budget stay apart.
2834///
2835/// A source's two readings are a plugin's word, so they are held to [`Metering`]'s contract
2836/// before either is believed, and a pair that breaks it is that source not metering — see
2837/// [`difference`].
2838pub(super) fn spent_between(
2839 before: &[Option<Metering>],
2840 after: &[Option<Metering>],
2841) -> Option<Spent> {
2842 let mut metered = false;
2843 let mut requests = 0_u64;
2844 let mut budgets: BTreeMap<(String, String), (u64, u64)> = BTreeMap::new();
2845 for (before, after) in before.iter().zip(after) {
2846 let (Some(before), Some(after)) = (before, after) else {
2847 continue;
2848 };
2849 let Some((sent, spent)) = difference(before, after) else {
2850 continue;
2851 };
2852 metered = true;
2853 requests = requests.saturating_add(sent);
2854 for (key, measured, modelled) in spent {
2855 let total = budgets.entry(key).or_default();
2856 total.0 = total.0.saturating_add(measured).saturating_add(modelled);
2857 total.1 = total.1.saturating_add(modelled);
2858 }
2859 }
2860 metered.then(|| Spent {
2861 requests,
2862 budgets: budgets
2863 .into_iter()
2864 .map(|((budget, unit), (amount, modelled))| BudgetSpent {
2865 budget,
2866 unit,
2867 amount,
2868 lower_bound: modelled > 0,
2869 })
2870 .collect(),
2871 })
2872}
2873
2874/// One budget's name and unit, and the measured and modelled amounts spent against it.
2875type BudgetDifference = ((String, String), u64, u64);
2876
2877/// What one source sent and spent between two of its readings, or `None` when the pair
2878/// cannot be a running total's.
2879///
2880/// A running total never falls, never drops a budget it has named, and names every budget
2881/// it keeps, once; a pair that breaks any of that is a source that reset its figures or
2882/// reported something else, and a difference taken over it would be a number that measures
2883/// nothing. Such a source is reported as not metering rather than as having spent a clamped
2884/// zero.
2885fn difference(before: &Metering, after: &Metering) -> Option<(u64, Vec<BudgetDifference>)> {
2886 let sent = after.requests.checked_sub(before.requests)?;
2887 for reading in [before, after] {
2888 let mut named = BTreeSet::new();
2889 for budget in &reading.budgets {
2890 if budget.budget.is_empty()
2891 || budget.unit.is_empty()
2892 || !named.insert((&budget.budget, &budget.unit))
2893 {
2894 return None;
2895 }
2896 }
2897 }
2898 let held = |from: &Metering, budget: &onetaskgraph_plugin_api::Metered| {
2899 from.budgets
2900 .iter()
2901 .find(|held| held.budget == budget.budget && held.unit == budget.unit)
2902 .cloned()
2903 };
2904 if before
2905 .budgets
2906 .iter()
2907 .any(|budget| held(after, budget).is_none())
2908 {
2909 return None;
2910 }
2911 let mut spent = Vec::with_capacity(after.budgets.len());
2912 for budget in &after.budgets {
2913 let earlier = held(before, budget);
2914 let measured = budget
2915 .measured
2916 .checked_sub(earlier.as_ref().map_or(0, |held| held.measured))?;
2917 let modelled = budget
2918 .modelled
2919 .checked_sub(earlier.as_ref().map_or(0, |held| held.modelled))?;
2920 spent.push((
2921 (budget.budget.clone(), budget.unit.clone()),
2922 measured,
2923 modelled,
2924 ));
2925 }
2926 Some((sent, spent))
2927}
2928
2929/// One page request against `source`, at the largest page it will serve.
2930fn request_for(source: &ResolvedSource, cursor: Option<Cursor>) -> PageRequest {
2931 PageRequest {
2932 cursor,
2933 limit: source.source().capabilities().max_page_size.max(1),
2934 }
2935}
2936
2937/// Whether writing this item would change what the destination already holds.
2938///
2939/// A free function over the state already read rather than a method that reads it again:
2940/// the same answer is wanted where the item is landed and where a repeat copy of a project
2941/// decides whether it settled, and a second read there is a second round trip for nothing.
2942fn changes(
2943 held: Option<&Prior>,
2944 item: &Planned,
2945 target: &NativeId,
2946 project: &Option<NativeId>,
2947 edges: &[DependencyEdge],
2948 delivers: &[TaskRef],
2949 destination: &SourceName,
2950) -> bool {
2951 let Some(held) = held else {
2952 return true;
2953 };
2954 let outgoing = outgoing(
2955 item,
2956 target.clone(),
2957 project.clone(),
2958 &recorded(item, Some(held)),
2959 delivers,
2960 Some(held),
2961 );
2962 !same(&held.item, &outgoing, destination) || !same_edges(&held.edges, edges)
2963}
2964
2965/// Remove one item this copy created, through the destination's own write interface.
2966async fn remove(
2967 destination: &ResolvedSource,
2968 kind: Level,
2969 id: &NativeId,
2970) -> Result<(), SourceError> {
2971 match kind {
2972 Level::Task => destination.source().delete_task(id).await,
2973 Level::Project => destination.source().delete_project(id).await,
2974 Level::Document => destination.source().delete_document(id).await,
2975 }
2976}
2977
2978/// Write one item back exactly as the destination held it before this copy.
2979async fn restore(
2980 destination: &ResolvedSource,
2981 id: &NativeId,
2982 prior: &Prior,
2983) -> Result<(), SourceError> {
2984 match &prior.item {
2985 Item::Task(task) => destination
2986 .source()
2987 .write_task(&ItemWrite {
2988 target: Some(id.clone()),
2989 item: (**task).clone(),
2990 depends_on: prior.edges.clone(),
2991 })
2992 .await
2993 .map(|_| ()),
2994 Item::Project(project) => destination
2995 .source()
2996 .write_project(&ItemWrite {
2997 target: Some(id.clone()),
2998 item: (**project).clone(),
2999 depends_on: prior.edges.clone(),
3000 })
3001 .await
3002 .map(|_| ()),
3003 Item::Document(document) => destination
3004 .source()
3005 .write_document(&ItemWrite {
3006 target: Some(id.clone()),
3007 item: (**document).clone(),
3008 depends_on: Vec::new(),
3009 })
3010 .await
3011 .map(|_| ()),
3012 }
3013}
3014
3015/// Hold `value` at `onetaskgraph.copies` on one item, through its source's narrow metadata
3016/// write — `false` when the source holds no such item.
3017async fn set_link(
3018 source: &ResolvedSource,
3019 level: Level,
3020 id: &NativeId,
3021 value: &Value,
3022) -> Result<bool, SourceError> {
3023 let key = MetadataKey::copies();
3024 Ok(match level {
3025 Level::Task => source
3026 .source()
3027 .set_task_metadata(id, &key, value)
3028 .await?
3029 .is_some(),
3030 Level::Project => source
3031 .source()
3032 .set_project_metadata(id, &key, value)
3033 .await?
3034 .is_some(),
3035 Level::Document => source
3036 .source()
3037 .set_document_metadata(id, &key, value)
3038 .await?
3039 .is_some(),
3040 })
3041}
3042
3043/// The metadata of either kind of item, to edit.
3044fn metadata_of(item: &mut Item) -> &mut BTreeMap<String, Value> {
3045 match item {
3046 Item::Task(task) => &mut task.metadata,
3047 Item::Project(project) => &mut project.metadata,
3048 Item::Document(document) => &mut document.metadata,
3049 }
3050}
3051
3052/// Refuse a document copy addressed to a source that declares it has none.
3053///
3054/// Read off the declaration rather than by asking, which is what "not asked" means: the
3055/// engine learned at the handshake that this source holds no documents, so it refuses
3056/// naming the source and its plugin instead of sending a read that would be refused there.
3057/// Applied at both ends of a copy — a source with no documents holds nothing to copy out,
3058/// and a destination with none has nowhere to put one.
3059fn documentary(source: &ResolvedSource) -> Result<(), EngineError> {
3060 if source.source().capabilities().documents.is_native() {
3061 return Ok(());
3062 }
3063 Err(EngineError::NoDocuments {
3064 name: source.name().to_string(),
3065 kind: source.kind().to_owned(),
3066 })
3067}
3068
3069/// One source failing while a copy was mid-flight.
3070fn refused(source: &ResolvedSource, error: SourceError) -> EngineError {
3071 EngineError::SourceRefused {
3072 name: source.name().to_string(),
3073 error,
3074 }
3075}
3076
3077/// Every forward edge at one item, walked to exhaustion one page at a time.
3078async fn forward_edges(
3079 source: &ResolvedSource,
3080 id: &NativeId,
3081 kind: Level,
3082) -> Result<Vec<DependencyEdge>, EngineError> {
3083 // A document has no edges to walk, and asking for them would mean asking a source for
3084 // a graph the contract says nothing may point into.
3085 if kind == Level::Document {
3086 return Ok(Vec::new());
3087 }
3088 let mut edges = Vec::new();
3089 let mut cursor: Option<Cursor> = None;
3090 loop {
3091 let asked = cursor.clone();
3092 let request = request_for(source, cursor);
3093 let page = match kind {
3094 Level::Task | Level::Document => {
3095 source
3096 .source()
3097 .task_dependencies(id, Direction::DependsOn, &request)
3098 .await
3099 }
3100 Level::Project => {
3101 source
3102 .source()
3103 .project_dependencies(id, Direction::DependsOn, &request)
3104 .await
3105 }
3106 }
3107 .map_err(|error| refused(source, error))?;
3108 fits(page.items.len(), request.limit).map_err(|error| refused(source, error))?;
3109 edges.extend(page.items);
3110 unrepeated(
3111 page.next.as_ref(),
3112 asked.as_ref(),
3113 "an item's dependencies were being read for a copy",
3114 )
3115 .map_err(|error| refused(source, error))?;
3116 match page.next {
3117 Some(next) => cursor = Some(next),
3118 None => return Ok(edges),
3119 }
3120 }
3121}
3122
3123/// The location string one record reports, when it reports a usable one.
3124///
3125/// Either variant's own `String`, and `None` for a record the source gave no location for
3126/// or gave an empty string for: there is nothing to look for in a document's content and
3127/// nothing to point a reader at.
3128fn located(location: Option<&Location>) -> Option<String> {
3129 let (Location::Path(held) | Location::Url(held)) = location?;
3130 (!held.is_empty()).then(|| held.clone())
3131}
3132
3133/// Record one candidate referent, when its source said where it is.
3134fn note(
3135 into: &mut Vec<Referent>,
3136 id: GlobalId,
3137 level: Level,
3138 location: Option<&Location>,
3139 metadata: &BTreeMap<String, Value>,
3140) {
3141 if let Some(location) = located(location) {
3142 into.push(Referent {
3143 id,
3144 origin: origin_of(metadata),
3145 level,
3146 location,
3147 });
3148 }
3149}
3150
3151/// What every location string this document names becomes, longest first.
3152///
3153/// Longest first because a shorter location may start where a longer one does — a
3154/// project's directory and a task's file under it — and the longer of the two is the
3155/// record that occurrence names.
3156fn table_for(referents: &[Referent], counterparts: &Counterparts) -> Vec<(String, Resolution)> {
3157 let mut table: Vec<(String, Resolution)> = Vec::new();
3158 for referent in referents {
3159 // Two referents reporting one location string: an occurrence of it cannot be
3160 // attributed to either, and a rewrite would be *confidently wrong* rather than
3161 // merely unhelpful. So neither is chosen and both occurrences are counted.
3162 if let Some(held) = table
3163 .iter_mut()
3164 .find(|(location, _)| location == &referent.location)
3165 {
3166 held.1 = Resolution::Ambiguous;
3167 continue;
3168 }
3169 table.push((referent.location.clone(), counterparts.resolve(referent)));
3170 }
3171 table.sort_by_key(|(location, _)| std::cmp::Reverse(location.len()));
3172 table
3173}
3174
3175/// Whether `content` holds `location` at least once, stopped on both sides.
3176fn holds(content: &str, location: &str) -> bool {
3177 (0..content.len()).any(|at| delimited_at(content, at, location))
3178}
3179
3180/// Whether `location` occurs at `at` **stopped on both sides** — by
3181/// [`stops_a_location`], or by the end of the content — rather than as part of a longer
3182/// location-like string.
3183///
3184/// A location string occurring inside a longer one is a different string naming a
3185/// different record: `/…/tasks/p/t.md` must not be rewritten inside `/…/tasks/p/t.md.bak`,
3186/// `https://example.invalid/1` must not be rewritten inside `https://example.invalid/12`,
3187/// and a project's location that is a directory prefix of a task's must not be rewritten
3188/// inside that task's. What deciding it this way costs is stated on [`stops_a_location`].
3189fn delimited_at(content: &str, at: usize, location: &str) -> bool {
3190 if !content.is_char_boundary(at) || !content[at..].starts_with(location) {
3191 return false;
3192 }
3193 let before = content[..at].chars().next_back();
3194 let after = content[at + location.len()..].chars().next();
3195 stops_a_location(before) && stops_a_location(after)
3196}
3197
3198/// Whether a character cannot continue a path or a link, so a location string beside one
3199/// ends there.
3200///
3201/// Stated as what *stops* a location rather than as what one may contain, because the
3202/// second list is unbounded — a path may hold very nearly any byte, and a URL more. Every
3203/// character not named here continues, which is what leaves the three cases above alone;
3204/// the end of the content counts as a stop. The set is what the artifact this exists for
3205/// really wraps a bare path in — a backtick in a table cell — plus the delimiters prose
3206/// and Markdown put next to one.
3207///
3208/// **Sentence punctuation is deliberately absent, and that is a stated cost rather than an
3209/// oversight.** `.`, `!`, `?` and `:` each equally *continue* a real location — `/…/t.md`
3210/// and `/…/t.md.bak` are two files, `…/1` and `…/1?q=2` two pages — so admitting them as
3211/// stops would rewrite one record's location into another's. The price is that a location
3212/// written bare at the end of a sentence is not recognised at all: its text is left
3213/// byte-for-byte and it is counted in neither figure, exactly as a reference to another
3214/// project's record is. That is the direction to be wrong in, because this edits the
3215/// content of somebody's document, where a confidently wrong rewrite is worse than one
3216/// that never happens.
3217fn stops_a_location(character: Option<char>) -> bool {
3218 match character {
3219 None => true,
3220 Some(character) => character.is_whitespace() || "`\"'()[]{}<>|,;".contains(character),
3221 }
3222}
3223
3224/// One document's content with every whole reference rewritten, and what that took.
3225///
3226/// A location string with no confident counterpart is left **byte-for-byte** as it was
3227/// rather than removed or guessed at, and so is every character of the content that is not
3228/// a rewritten reference. Nothing is added and nothing is reformatted.
3229fn substitute(content: &str, table: &[(String, Resolution)]) -> (String, Counted) {
3230 let mut written = String::with_capacity(content.len());
3231 let mut counts = Counted::default();
3232 let mut at = 0;
3233 while at < content.len() {
3234 if let Some((location, resolution)) = table
3235 .iter()
3236 .find(|(location, _)| delimited_at(content, at, location))
3237 {
3238 match resolution {
3239 Resolution::Rewrite(there) => {
3240 written.push_str(there);
3241 counts.rewritten += 1;
3242 }
3243 // Both left-alone outcomes count as unresolved on the branch that counts
3244 // them, which is what really holds `ambiguous` at or below `unresolved`.
3245 Resolution::NoCounterpart => {
3246 written.push_str(location);
3247 counts.unresolved += 1;
3248 }
3249 Resolution::Ambiguous => {
3250 written.push_str(location);
3251 counts.unresolved += 1;
3252 counts.ambiguous += 1;
3253 }
3254 }
3255 at += location.len();
3256 continue;
3257 }
3258 let character = content[at..]
3259 .chars()
3260 .next()
3261 .expect("a character at a boundary this walk only ever lands on");
3262 written.push(character);
3263 at += character.len_utf8();
3264 }
3265 (written, counts)
3266}
3267
3268/// The origin one item records, when it records a usable one.
3269fn origin_of(metadata: &BTreeMap<String, Value>) -> Option<GlobalId> {
3270 metadata
3271 .get(GlobalId::ORIGIN_KEY)?
3272 .as_str()?
3273 .parse::<GlobalId>()
3274 .ok()
3275}
3276
3277/// The destination item one item's `onetaskgraph.copies` entry names at `destination`, when
3278/// it holds a usable one.
3279///
3280/// An entry that is not a qualified id, or is one naming another source than the key it is
3281/// filed under, is no link at all: the copy goes on by the rules that follow and rewrites
3282/// it to what they find.
3283fn link_of(metadata: &BTreeMap<String, Value>, destination: &SourceName) -> Option<GlobalId> {
3284 let linked = metadata
3285 .get(MetadataKey::COPIES_KEY)?
3286 .as_object()?
3287 .get(destination.as_str())?
3288 .as_str()?
3289 .parse::<GlobalId>()
3290 .ok()?;
3291 (&linked.source == destination).then_some(linked)
3292}
3293
3294/// The entries of a `onetaskgraph.copies` value that are links: a source name, and a
3295/// qualified id of that same source.
3296fn well_formed(value: Option<&Value>) -> serde_json::Map<String, Value> {
3297 let Some(Value::Object(held)) = value else {
3298 return serde_json::Map::new();
3299 };
3300 held.iter()
3301 .filter(|(destination, linked)| {
3302 linked
3303 .as_str()
3304 .and_then(|linked| linked.parse::<GlobalId>().ok())
3305 .is_some_and(|linked| linked.source.as_str() == destination.as_str())
3306 })
3307 .map(|(destination, linked)| (destination.clone(), linked.clone()))
3308 .collect()
3309}
3310
3311/// Why a value is not one a `onetaskgraph.copies` entry may hold, or nothing when it is: a
3312/// JSON object mapping each destination source name to a qualified id of that source.
3313///
3314/// What the stdio boundary holds an incoming value of that key to, so a plugin is never
3315/// handed a link it could not answer a later copy with.
3316pub(crate) fn malformed_links(value: &Value) -> Option<String> {
3317 let Value::Object(held) = value else {
3318 return Some(format!(
3319 "{} holds an object of destination source names to qualified ids, not {value}",
3320 MetadataKey::COPIES_KEY
3321 ));
3322 };
3323 let wrong = held.len() - well_formed(Some(value)).len();
3324 (wrong > 0).then(|| {
3325 format!(
3326 "{} holds an object of destination source names to qualified ids of that source, \
3327 and {wrong} of its entries is not one: {value}",
3328 MetadataKey::COPIES_KEY
3329 )
3330 })
3331}
3332
3333/// The title and metadata of either kind of item.
3334fn described(item: &Item) -> (&str, &BTreeMap<String, Value>) {
3335 match item {
3336 Item::Task(task) => (&task.title, &task.metadata),
3337 Item::Project(project) => (&project.title, &project.metadata),
3338 Item::Document(document) => (&document.title, &document.metadata),
3339 }
3340}
3341
3342/// The item as the destination should hold it.
3343///
3344/// `url`, `location`, `created_at` and `updated_at` are the destination's own and are
3345/// never written — where the *source* holds an item says nothing about where the
3346/// destination does, which is why a copied document does not arrive claiming the path or
3347/// the link its source reported. The two reserved keys this product encodes typed fields
3348/// under are removed, because those fields travel as themselves — leaving the encoding
3349/// beside them would have the destination hold one thing twice, and disagree with itself
3350/// the moment one changed.
3351///
3352/// A task's `delivers` is `delivers`, already resolved against the destination, and its
3353/// `delivered_by` is the one the destination holds — never the source's, and empty for an
3354/// item this copy creates: that list is the store's to keep, at the destination. The
3355/// `onetaskgraph.copies` link is kept on the same terms: it says where the *destination*
3356/// item was copied to, which is the destination's to hold, never the source's to send.
3357fn outgoing(
3358 item: &Planned,
3359 id: NativeId,
3360 project: Option<NativeId>,
3361 origin: &Origin,
3362 delivers: &[TaskRef],
3363 held: Option<&Prior>,
3364) -> Item {
3365 let own = held.and_then(|held| described(&held.item).1.get(MetadataKey::COPIES_KEY));
3366 let carried = |metadata: &BTreeMap<String, Value>| carried(metadata, origin, own);
3367 match &item.item {
3368 Item::Task(task) => Item::Task(Box::new(Task {
3369 id,
3370 url: None,
3371 location: None,
3372 created_at: None,
3373 updated_at: None,
3374 project,
3375 metadata: carried(&task.metadata),
3376 delivers: delivers.to_vec(),
3377 delivered_by: match held.map(|held| &held.item) {
3378 Some(Item::Task(held)) => held.delivered_by.clone(),
3379 _ => Vec::new(),
3380 },
3381 ..(**task).clone()
3382 })),
3383 Item::Project(project) => Item::Project(Box::new(Project {
3384 id,
3385 url: None,
3386 location: None,
3387 created_at: None,
3388 updated_at: None,
3389 metadata: carried(&project.metadata),
3390 ..(**project).clone()
3391 })),
3392 Item::Document(document) => Item::Document(Box::new(Document {
3393 id,
3394 url: None,
3395 location: None,
3396 created_at: None,
3397 updated_at: None,
3398 project,
3399 metadata: carried(&document.metadata),
3400 ..(**document).clone()
3401 })),
3402 }
3403}
3404
3405/// The id a created item is offered to the destination under.
3406///
3407/// A task whose id is `<its project's id>/<rest>` is offered as `<rest>` under the
3408/// destination project it is filed in, so a destination whose ids are paths files it in that
3409/// project rather than the source's. Any other id, and every update, is left as it is. The
3410/// `a_project_copied_*` journeys in `crates/onetaskgraph/tests/e2e/copy.rs` hold the Markdown
3411/// plugin's id shape to this rule.
3412fn created_id(item: &Item, filed: Option<&NativeId>) -> NativeId {
3413 if let (Item::Task(task), Some(filed)) = (item, filed)
3414 && let Some(own) = &task.project
3415 && let Some(rest) = task
3416 .id
3417 .as_str()
3418 .strip_prefix(own.as_str())
3419 .and_then(|rest| rest.strip_prefix('/'))
3420 && !rest.is_empty()
3421 {
3422 return NativeId(format!("{}/{rest}", filed.as_str()));
3423 }
3424 item.id().clone()
3425}
3426
3427/// The metadata a copy carries: the caller's own keys untouched, the origin settled, and the
3428/// destination's own `onetaskgraph.copies` link — `own` — kept as it holds it.
3429///
3430/// The key is removed before it is settled rather than overwritten, because the item being
3431/// copied carries an origin of its own and [`Origin::Keeps`] must not let it through. Its
3432/// link is removed for the same reason: it names where *that* item was copied to, and on
3433/// the destination it would name the destination item's own counterpart as itself.
3434fn carried(
3435 metadata: &BTreeMap<String, Value>,
3436 origin: &Origin,
3437 own: Option<&Value>,
3438) -> BTreeMap<String, Value> {
3439 let mut carried = metadata.clone();
3440 carried.remove(Repository::METADATA_KEY);
3441 carried.remove(DependencyEdge::RECORDED_KEY);
3442 carried.remove(TaskRef::DELIVERS_KEY);
3443 carried.remove(TaskRef::DELIVERED_BY_KEY);
3444 carried.remove(GlobalId::ORIGIN_KEY);
3445 carried.remove(MetadataKey::COPIES_KEY);
3446 if let Some(own) = own {
3447 carried.insert(MetadataKey::COPIES_KEY.to_owned(), own.clone());
3448 }
3449 let held = match origin {
3450 Origin::Records(id) => Some(Value::String(id.to_string())),
3451 Origin::Keeps(held) => held.clone(),
3452 };
3453 if let Some(held) = held {
3454 carried.insert(GlobalId::ORIGIN_KEY.to_owned(), held);
3455 }
3456 carried
3457}
3458
3459/// What one landed item records at [`GlobalId::ORIGIN_KEY`].
3460enum Origin {
3461 /// The qualified id this item was copied from, as the id type rather than as its
3462 /// spelling: the key holds a [`GlobalId`] and nothing else may be recorded there.
3463 Records(GlobalId),
3464 /// Whatever the destination already holds there — `None` when it holds nothing, which
3465 /// is written as the key being absent rather than as a null.
3466 ///
3467 /// A [`Value`] and not a [`GlobalId`], because this variant does not interpret what it
3468 /// carries: it is the destination's own metadata entry, held for the length of one
3469 /// write and put back exactly as it was read. Parsing it would turn a value a
3470 /// destination holds and this engine cannot read into a value this engine deletes,
3471 /// which is the opposite of what keeping it means.
3472 Keeps(Option<Value>),
3473}
3474
3475/// Which of the two a copy of this item does.
3476///
3477/// A copy that reached its target by rule 1 is a copy-back: the item being copied names
3478/// the destination item, so the destination is the *original* and the id being copied
3479/// belongs to the copy that came out of it. Recording that id there would overwrite the
3480/// original's own provenance — and with it the correspondence every later copy from the
3481/// source it was authored in depends on. That copy would then match nothing and create a
3482/// second item beside the one it meant to update, which is the whole failure: nothing is
3483/// reported, and whoever reads that board now has two. So a copy-back leaves the
3484/// destination's origin exactly as the destination holds it, absent included, and every
3485/// other copy records the id it was copied from.
3486fn recorded(item: &Planned, held: Option<&Prior>) -> Origin {
3487 if let Target::Update {
3488 found: Found::Origin,
3489 ..
3490 } = &item.target
3491 {
3492 return Origin::Keeps(
3493 held.and_then(|held| described(&held.item).1.get(GlobalId::ORIGIN_KEY).cloned()),
3494 );
3495 }
3496 Origin::Records(item.source.clone())
3497}
3498
3499/// Whether the destination already reads exactly as this copy would leave it.
3500///
3501/// The destination's own `url` and timestamps are excluded because a copy never writes
3502/// them, so a difference there is not one this copy would close.
3503fn same(held: &Item, outgoing: &Item, destination: &SourceName) -> bool {
3504 match (held, outgoing) {
3505 (Item::Task(held), Item::Task(outgoing)) => {
3506 // Qualified before they are compared, so `T-1` and `folder:T-1` at the destination
3507 // `folder` are the one entry they are.
3508 targets(&held.delivers, destination) == targets(&outgoing.delivers, destination)
3509 && targets(&held.delivered_by, destination)
3510 == targets(&outgoing.delivered_by, destination)
3511 && held.title == outgoing.title
3512 && held.content == outgoing.content
3513 && held.status == outgoing.status
3514 && held.priority == outgoing.priority
3515 && held.labels == outgoing.labels
3516 && held.project == outgoing.project
3517 && held.metadata == outgoing.metadata
3518 && held.repositories == outgoing.repositories
3519 }
3520 (Item::Project(held), Item::Project(outgoing)) => {
3521 held.title == outgoing.title
3522 && held.content == outgoing.content
3523 && held.status == outgoing.status
3524 && held.labels == outgoing.labels
3525 && held.metadata == outgoing.metadata
3526 && held.repositories == outgoing.repositories
3527 }
3528 // No status, because a document has none; no edges, because it is in no graph.
3529 (Item::Document(held), Item::Document(outgoing)) => {
3530 held.title == outgoing.title
3531 && held.content == outgoing.content
3532 && held.labels == outgoing.labels
3533 && held.project == outgoing.project
3534 && held.metadata == outgoing.metadata
3535 && held.repositories == outgoing.repositories
3536 }
3537 _ => false,
3538 }
3539}
3540
3541/// Whether the destination's forward edges already say what this copy would write.
3542fn same_edges(held: &[DependencyEdge], outgoing: &[DependencyEdge]) -> bool {
3543 let ends = |edges: &[DependencyEdge]| {
3544 let mut ends: Vec<(String, ItemKind, DependencyKind)> = edges
3545 .iter()
3546 .map(|edge| (edge.to.id().to_owned(), edge.to.kind, edge.kind))
3547 .collect();
3548 ends.sort_by(|left, right| left.0.cmp(&right.0));
3549 ends
3550 };
3551 ends(held) == ends(outgoing)
3552}
3553
3554/// Each read edge as the destination should record it, or `None` when its far end is a
3555/// member of this copy whose destination id is not known yet.
3556fn mapped_edges(
3557 edges: &[DependencyEdge],
3558 origin: &SourceName,
3559 destination: &ResolvedSource,
3560 copied: &[GlobalId],
3561 written: &BTreeMap<String, NativeId>,
3562) -> Vec<Option<DependencyEdge>> {
3563 edges
3564 .iter()
3565 .map(|edge| {
3566 let far = GlobalId::new(origin.clone(), NativeId(edge.to.id().to_owned()));
3567 let id = if let Some(native) = names(&edge.to, destination.name()) {
3568 // A far end already qualified to the destination's own source is that
3569 // source's own item, so it is written the way that source names its own:
3570 // unqualified. Leaving it qualified would have the destination hold an
3571 // edge into itself written as if it left, which is the one spelling the
3572 // reserved key exists to keep for edges that really do.
3573 Some(native)
3574 } else if !edge.to.is_qualified() && copied.contains(&far) {
3575 // A member of this copy, inside one source included: a copy there lands
3576 // beside the item it was read from, so the edge names the copy rather than
3577 // the original the copy's own project does not hold.
3578 written.get(&far.to_string()).map(|native| native.0.clone())
3579 } else if edge.to.is_qualified() || origin == destination.name() {
3580 // Already naming a source of its own, or a copy inside one source where
3581 // the far end's own id is the destination's id.
3582 Some(edge.to.id().to_owned())
3583 } else {
3584 Some(far.to_string())
3585 }?;
3586 DependencyEndpoint::new(id, edge.to.kind)
3587 .ok()
3588 .map(|to| DependencyEdge {
3589 from: edge.from.clone(),
3590 to,
3591 kind: edge.kind,
3592 })
3593 })
3594 .collect()
3595}
3596
3597/// Each `delivers` entry of one planned task as the destination should hold it, or `None`
3598/// where it names a member of this copy whose destination id is not known yet. Empty for
3599/// anything but a task.
3600fn delivers_of(
3601 item: &Planned,
3602 destination: &ResolvedSource,
3603 copied: &[GlobalId],
3604 written: &BTreeMap<String, NativeId>,
3605) -> Vec<Option<TaskRef>> {
3606 match &item.item {
3607 Item::Task(task) => mapped_delivers(
3608 &task.delivers,
3609 &item.source.source,
3610 destination,
3611 copied,
3612 written,
3613 ),
3614 Item::Project(_) | Item::Document(_) => Vec::new(),
3615 }
3616}
3617
3618/// Each `delivers` entry as the destination should hold it, or `None` when it names a member
3619/// of this copy whose destination id is not known yet.
3620///
3621/// An entry naming a member of the copied set becomes that member's own id at the
3622/// destination — bare, because the member is the destination's, unless the id holds a colon
3623/// a bare entry would be misread by. Every other entry is carried through qualified, so a
3624/// bare one read at `origin` goes on naming a task of `origin`.
3625fn mapped_delivers(
3626 entries: &[TaskRef],
3627 origin: &SourceName,
3628 destination: &ResolvedSource,
3629 copied: &[GlobalId],
3630 written: &BTreeMap<String, NativeId>,
3631) -> Vec<Option<TaskRef>> {
3632 entries
3633 .iter()
3634 .map(|entry| {
3635 let qualified = entry.in_source(origin);
3636 let Ok(far) = qualified.as_str().parse::<GlobalId>() else {
3637 return Some(qualified);
3638 };
3639 if !copied.contains(&far) {
3640 return Some(qualified);
3641 }
3642 let native = written.get(&far.to_string())?;
3643 Some(if native.as_str().contains(':') {
3644 TaskRef::qualified(destination.name(), native)
3645 } else {
3646 TaskRef::new(native.as_str())
3647 .unwrap_or_else(|_| TaskRef::qualified(destination.name(), native))
3648 })
3649 })
3650 .collect()
3651}
3652
3653/// How many of one task's `delivers` entries name a member of the copied set.
3654fn members_named(entries: &[TaskRef], origin: &SourceName, copied: &[GlobalId]) -> u64 {
3655 let named = entries
3656 .iter()
3657 .filter(|entry| {
3658 entry
3659 .in_source(origin)
3660 .as_str()
3661 .parse::<GlobalId>()
3662 .is_ok_and(|far| copied.contains(&far))
3663 })
3664 .count();
3665 u64::try_from(named).unwrap_or(u64::MAX)
3666}
3667
3668/// The entries that could be resolved, which is every one of them on the second pass.
3669fn resolved_entries(entries: &[Option<TaskRef>]) -> Vec<TaskRef> {
3670 entries.iter().flatten().cloned().collect()
3671}
3672
3673/// The native id a qualified endpoint names at `destination`, when it names one there.
3674fn names(endpoint: &DependencyEndpoint, destination: &SourceName) -> Option<String> {
3675 if !endpoint.is_qualified() {
3676 return None;
3677 }
3678 let id: GlobalId = endpoint.id().parse().ok()?;
3679 (&id.source == destination).then_some(id.native.0)
3680}
3681
3682/// The edges that could be resolved, which is every one of them on the second pass.
3683fn resolved(edges: &[Option<DependencyEdge>]) -> Vec<DependencyEdge> {
3684 edges.iter().flatten().cloned().collect()
3685}
3686
3687/// Refuse an item of a member copy whose edge names a member that copy does not carry
3688/// and whose destination id nothing records.
3689///
3690/// `unrecorded` is empty for every other copy, which is what makes this a no-op there. An
3691/// edge already naming a source of its own, and every edge of a copy inside one source,
3692/// is written as it was read by [`mapped_edges`], so neither can need a recorded origin.
3693fn unrecorded_far_end(
3694 item: &Planned,
3695 destination: &ResolvedSource,
3696 unrecorded: &[GlobalId],
3697) -> Result<(), EngineError> {
3698 if &item.source.source == destination.name() {
3699 return Ok(());
3700 }
3701 for edge in &item.edges {
3702 if edge.to.is_qualified() {
3703 continue;
3704 }
3705 let far = GlobalId::new(
3706 item.source.source.clone(),
3707 NativeId(edge.to.id().to_owned()),
3708 );
3709 if unrecorded.contains(&far) {
3710 return Err(EngineError::UnrecordedMember {
3711 item: item.source.clone(),
3712 member: far,
3713 destination: destination.name().clone(),
3714 });
3715 }
3716 }
3717 Ok(())
3718}