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