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