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