1use std::path::Path;
9
10use omgbase_format::BlockKind;
11use omgbase_format::text::{normalize_text, normalize_visible_text};
12use omgbase_reconcile::Config;
13use omgbase_search::EmbeddingProvider;
14use omgbase_store::docs_ops::diffs_json;
15use omgbase_store::mutate_kernel::{At, Op, Parent, To};
16use omgbase_store::{
17 ApplyOrigin, ApplyRequest, ApplyResult, DocOpContext, DocStore, Expect, FsDocStore, Opset,
18 QueryVector, Store,
19};
20use omgbase_sync::{RealFileSystem, RepoRow};
21use rusqlite::{OptionalExtension, params};
22use serde_json::{Map, Value as Json, json};
23
24use crate::error::{Result, SurfaceError};
25use crate::graph::{GraphArgs, graph_neighborhood};
26use crate::history;
27use crate::links;
28use crate::query::{QueryOptions, query};
29use crate::read::{self, Resolution, ResolvedRef};
30use crate::reference::QUERY_SYNTAX;
31
32pub const ACTOR: &str = "agent:mcp";
34
35#[derive(Clone, Debug, PartialEq)]
37pub struct ToolSpec {
38 pub name: &'static str,
39 pub description: &'static str,
40 pub input_schema: Json,
42}
43
44#[derive(Clone, Debug, PartialEq)]
47pub struct ToolOutcome {
48 pub body: Json,
49 pub is_error: bool,
50}
51
52enum WriteTarget {
54 Derived,
57 Fixed(Box<dyn DocStore>),
60}
61
62pub struct Surface {
64 store: Store,
65 default_repo: String,
66 provider: Option<Box<dyn EmbeddingProvider>>,
67 on_mutation: Option<Box<dyn FnMut()>>,
68 clock: Box<dyn FnMut() -> String>,
69 writes: WriteTarget,
70 config: Config,
71 actor: String,
73}
74
75fn bad_args(msg: impl Into<String>) -> SurfaceError {
78 SurfaceError::filter_invalid(msg, "arguments")
79}
80
81fn arg_str<'a>(args: &'a Json, key: &str) -> Option<&'a str> {
82 args.get(key).and_then(Json::as_str)
83}
84
85fn arg_string(args: &Json, key: &str) -> Result<String> {
86 arg_str(args, key)
87 .map(str::to_owned)
88 .ok_or_else(|| bad_args(format!("`{key}` must be a string")))
89}
90
91fn arg_i64(args: &Json, key: &str) -> Result<Option<i64>> {
92 match args.get(key) {
93 None | Some(Json::Null) => Ok(None),
94 Some(v) => v
95 .as_f64()
96 .filter(|n| n.fract() == 0.0)
97 .map(|n| Some(n as i64))
98 .ok_or_else(|| bad_args(format!("`{key}` must be an integer"))),
99 }
100}
101
102fn arg_usize(args: &Json, key: &str) -> Result<Option<usize>> {
103 Ok(arg_i64(args, key)?.map(|n| usize::try_from(n).unwrap_or(0)))
104}
105
106fn arg_bool(args: &Json, key: &str) -> Result<Option<bool>> {
107 match args.get(key) {
108 None | Some(Json::Null) => Ok(None),
109 Some(Json::Bool(b)) => Ok(Some(*b)),
110 Some(_) => Err(bad_args(format!("`{key}` must be a boolean"))),
111 }
112}
113
114fn arg_strings(args: &Json, key: &str) -> Result<Vec<String>> {
115 let Some(arr) = args.get(key).and_then(Json::as_array) else {
116 return Err(bad_args(format!("`{key}` must be an array of strings")));
117 };
118 arr.iter()
119 .map(|v| {
120 v.as_str()
121 .map(str::to_owned)
122 .ok_or_else(|| bad_args(format!("`{key}` must be an array of strings")))
123 })
124 .collect()
125}
126
127fn arg_object<'a>(args: &'a Json, key: &str) -> Result<Option<&'a Map<String, Json>>> {
128 match args.get(key) {
129 None | Some(Json::Null) => Ok(None),
130 Some(Json::Object(o)) => Ok(Some(o)),
131 Some(_) => Err(bad_args(format!("`{key}` must be an object"))),
132 }
133}
134
135fn hex_field(e: &Map<String, Json>, k: &str) -> Option<String> {
136 e.get(k).and_then(Json::as_str).map(str::to_owned)
137}
138
139fn arg_expect(args: &Json) -> Result<Option<Expect>> {
142 Ok(arg_object(args, "expect")?.map(|e| Expect {
143 content_hash: hex_field(e, "content_hash"),
144 parent_children_hash: hex_field(e, "parent_children_hash"),
145 }))
146}
147
148fn arg_parent_expect(args: &Json) -> Result<Option<Expect>> {
153 Ok(arg_object(args, "expect")?.map(|e| Expect {
154 content_hash: None,
155 parent_children_hash: hex_field(e, "parent_children_hash"),
156 }))
157}
158
159fn resolution_arg(args: &Json, default: Resolution) -> Result<Resolution> {
160 match arg_str(args, "resolution") {
161 None => Ok(default),
162 Some(s) => Resolution::parse(s).ok_or_else(|| {
163 bad_args("`resolution` must be one of skeleton, outline, text, raw, full")
164 }),
165 }
166}
167
168fn merge(mut a: Map<String, Json>, b: Json) -> Json {
170 if let Json::Object(o) = b {
171 for (k, v) in o {
172 a.insert(k, v);
173 }
174 }
175 Json::Object(a)
176}
177
178fn schema(props: &[(&str, Json)], required: &[&str], repo: bool) -> Json {
179 let mut p = Map::new();
180 for (k, v) in props {
181 p.insert((*k).to_owned(), v.clone());
182 }
183 if repo {
184 p.insert("repo".to_owned(), json!({ "type": "string" }));
185 }
186 json!({ "type": "object", "properties": p, "required": required })
187}
188
189fn s() -> Json {
190 json!({ "type": "string" })
191}
192fn i() -> Json {
193 json!({ "type": "integer" })
194}
195fn b() -> Json {
196 json!({ "type": "boolean" })
197}
198fn strings() -> Json {
199 json!({ "type": "array", "items": { "type": "string" } })
200}
201fn obj() -> Json {
202 json!({ "type": "object" })
203}
204fn nullable_string() -> Json {
205 json!({ "type": ["string", "null"] })
206}
207fn resolution_schema() -> Json {
208 json!({ "type": "string", "enum": ["skeleton", "outline", "text", "raw", "full"] })
209}
210fn at_schema() -> Json {
211 json!({ "oneOf": [
212 { "const": "start" }, { "const": "end" },
213 { "type": "object", "properties": { "before": { "type": "string" } }, "required": ["before"] },
214 { "type": "object", "properties": { "after": { "type": "string" } }, "required": ["after"] }
215 ] })
216}
217fn expect_schema() -> Json {
218 json!({ "type": "object", "properties": { "content_hash": { "type": "string" }, "parent_children_hash": { "type": "string" } } })
219}
220fn parent_expect_schema() -> Json {
225 json!({ "type": "object", "properties": { "parent_children_hash": { "type": "string" } } })
226}
227
228#[must_use]
230pub fn tools() -> Vec<ToolSpec> {
231 let dp = |extra: &[(&str, Json)], required: &[&str]| {
232 let mut props = vec![("doc", s()), ("path", s())];
233 props.extend(extra.iter().cloned());
234 schema(&props, required, true)
235 };
236 vec![
237 ToolSpec {
238 name: "docs_outline",
239 description: "A document's compact indented outline (id, type, label per line; § marks headings). Args take a doc id or path.",
240 input_schema: dp(
241 &[
242 (
243 "resolution",
244 json!({ "type": "string", "enum": ["skeleton", "outline"] }),
245 ),
246 ("depth", i()),
247 ("budget_tokens", i()),
248 ],
249 &[],
250 ),
251 },
252 ToolSpec {
253 name: "docs_read",
254 description: "Read a whole document: verbatim `content`, `properties` grouped by source, `path`/`docId`/`rev`; include_ids adds `ids`, `hashes` (CAS tokens) and `parents`.",
255 input_schema: dp(&[("include_ids", b())], &[]),
256 },
257 ToolSpec {
258 name: "docs_get_many",
259 description: "Batch docs_read over `docs` (ids or paths): `{ items, errors, truncated }`, capped at 100 refs, optional token budget.",
260 input_schema: schema(
261 &[
262 ("docs", strings()),
263 ("include_ids", b()),
264 ("budget_tokens", i()),
265 ],
266 &["docs"],
267 true,
268 ),
269 },
270 ToolSpec {
271 name: "nodes_get",
272 description: "Hydrate one block subtree at a resolution (skeleton|outline|text|raw|full); the owning doc is inferred from `id` when `doc`/`path` are omitted.",
273 input_schema: dp(&[("id", s()), ("resolution", resolution_schema())], &["id"]),
274 },
275 ToolSpec {
276 name: "nodes_get_many",
277 description: "Fetch up to 100 blocks by id in request order with budget truncation; `doc`/`path` is an optional scope. Returns `nodes`, `truncated`, `unresolved`.",
278 input_schema: dp(
279 &[
280 ("ids", strings()),
281 ("resolution", resolution_schema()),
282 ("budget_tokens", i()),
283 ],
284 &["ids"],
285 ),
286 },
287 ToolSpec {
288 name: "read_ref",
289 description: "Read any ref — a document (id or path) or a block (`b_` id, or an `n_` node id) — classified as `{ kind: \"document\", … }` or `{ kind: \"block\", … }` (block `resolution` defaults to raw).",
290 input_schema: schema(
291 &[("ref", s()), ("resolution", resolution_schema())],
292 &["ref"],
293 true,
294 ),
295 },
296 ToolSpec {
297 name: "docs_tree",
298 description: "The directory-aware shape of a repo: live docs under `path` collapsed at `depth` segments into dir/doc entries with totals, ordered by path and paged.",
299 input_schema: schema(
300 &[
301 ("path", s()),
302 ("depth", i()),
303 ("limit", i()),
304 ("cursor", nullable_string()),
305 ("budget_tokens", i()),
306 ],
307 &[],
308 true,
309 ),
310 },
311 ToolSpec {
312 name: "docs_list",
313 description: "Enumerate live documents as a page `{ items: [{ path, blocks, ts }], truncated, cursor }`, ordered by path; `path_glob` is a LIKE match (`*` matches across `/`).",
314 input_schema: schema(
315 &[
316 ("path_glob", s()),
317 ("limit", i()),
318 ("cursor", nullable_string()),
319 ("budget_tokens", i()),
320 ],
321 &[],
322 true,
323 ),
324 },
325 ToolSpec {
326 name: "query_syntax",
327 description: "The OQX syntax reference for the `query` tool.",
328 input_schema: schema(&[], &[], false),
329 },
330 ToolSpec {
331 name: "query",
332 description: "Run one OQX query (`select … from docs|blocks|nodes|edges where … follow … order by … limit N`). Returns lean hits `{ id, path, …projections }` with `truncated` + `cursor`, or a `count`/`exists`/`none` scalar, or `values`. See query_syntax.",
333 input_schema: schema(
334 &[
335 ("query", s()),
336 ("limit", i()),
337 ("cursor", nullable_string()),
338 ],
339 &["query"],
340 true,
341 ),
342 },
343 ToolSpec {
344 name: "graph",
345 description: "The bounded neighborhood around root documents in one call — `{ documents, edges, frontier }` — compiled to an OQX `follow doc.out`/`doc.in` walk.",
346 input_schema: schema(
347 &[
348 ("roots", strings()),
349 ("degrees", i()),
350 (
351 "direction",
352 json!({ "type": "string", "enum": ["in", "out", "both"] }),
353 ),
354 ("predicate", s()),
355 ("select", strings()),
356 ("max_documents", i()),
357 ],
358 &["roots"],
359 true,
360 ),
361 },
362 ToolSpec {
363 name: "text_search",
364 description: "Full-text (FTS5, bm25-ranked) search over block text.",
365 input_schema: schema(&[("q", s()), ("limit", i())], &["q"], true),
366 },
367 ToolSpec {
368 name: "resolve",
369 description: "Resolve a name/title/phrase to the blocks it refers to: ranked `{ id, locator, preview, evidence }` (FTS, fused with the vector ranking when a provider exists).",
370 input_schema: schema(&[("query", s()), ("limit", i())], &["query"], true),
371 },
372 ToolSpec {
373 name: "apply",
374 description: "Apply a changeset of kernel ops (insert/update/move/remove/split/merge) atomically; `dry_run` previews diffs.",
375 input_schema: schema(
376 &[
377 (
378 "ops",
379 json!({ "type": "array", "items": { "type": "object" } }),
380 ),
381 ("reason", s()),
382 ("dry_run", b()),
383 ],
384 &["ops"],
385 true,
386 ),
387 },
388 ToolSpec {
389 name: "blocks_insert",
390 description: "Insert blocks parsed from `markdown` under `to` (a block ref, or a document ref for its top level) at `at` (end|start|{before|after}). Optional `expect.parent_children_hash` is the destination-parent CAS: the parent's CURRENT direct child ids (the document's top-level ids for a document `to`) joined by `,` and sha256-hexed — compute it from docs_read include_ids (`ids` filtered by `parents`) — and the insert fails `stale_expectation` (with `data.current.parent_children_hash`) if the siblings changed under you; `content_hash` has no meaning here (there is no target block) and is not accepted.",
391 input_schema: schema(
392 &[
393 ("to", s()),
394 ("markdown", s()),
395 ("at", at_schema()),
396 ("expect", parent_expect_schema()),
397 ("dry_run", b()),
398 ],
399 &["to", "markdown"],
400 true,
401 ),
402 },
403 ToolSpec {
404 name: "blocks_update",
405 description: "Replace a block's markdown and/or set attrs (`checked` folds into attrs) with CAS pinned server-side when `expect` is omitted. Returns `id`, `ids` and the apply result.",
406 input_schema: schema(
407 &[
408 ("block", s()),
409 ("markdown", s()),
410 ("checked", b()),
411 ("attrs", obj()),
412 ("expect", expect_schema()),
413 ("dry_run", b()),
414 ],
415 &["block"],
416 true,
417 ),
418 },
419 ToolSpec {
420 name: "blocks_move",
421 description: "Move blocks under a new parent at a position; `to` is a block ref or the blocks' own document. Optional `expect.parent_children_hash` is the destination-parent CAS (as blocks_insert: the destination parent's CURRENT direct child ids joined by `,`, sha256 hex), checked once for the whole run before anything moves — `stale_expectation` with `data.current.parent_children_hash` if the destination's children changed; `content_hash` is meaningless here and not accepted.",
422 input_schema: schema(
423 &[
424 ("blocks", strings()),
425 ("to", s()),
426 ("at", at_schema()),
427 ("expect", parent_expect_schema()),
428 ("dry_run", b()),
429 ],
430 &["blocks", "to"],
431 true,
432 ),
433 },
434 ToolSpec {
435 name: "blocks_remove",
436 description: "Remove blocks (and their subtrees).",
437 input_schema: schema(
438 &[("blocks", strings()), ("dry_run", b())],
439 &["blocks"],
440 true,
441 ),
442 },
443 ToolSpec {
444 name: "blocks_split",
445 description: "Split a block at UTF-8 byte offsets; CAS pinned server-side.",
446 input_schema: schema(
447 &[
448 ("block", s()),
449 (
450 "at",
451 json!({ "type": "array", "items": { "type": "integer" } }),
452 ),
453 ("dry_run", b()),
454 ],
455 &["block", "at"],
456 true,
457 ),
458 },
459 ToolSpec {
460 name: "blocks_merge",
461 description: "Merge adjacent blocks into the first, joined by `separator`.",
462 input_schema: schema(
463 &[("blocks", strings()), ("separator", s()), ("dry_run", b())],
464 &["blocks"],
465 true,
466 ),
467 },
468 ToolSpec {
469 name: "tasks_complete",
470 description: "Check (or uncheck with checked:false) task blocks.",
471 input_schema: schema(
472 &[("blocks", strings()), ("checked", b()), ("dry_run", b())],
473 &["blocks"],
474 true,
475 ),
476 },
477 ToolSpec {
478 name: "node_set",
479 description: "Set one editable property of a projected node (a link's name/value, a task's checked).",
480 input_schema: schema(
481 &[
482 ("node", s()),
483 ("prop", s()),
484 ("value", s()),
485 ("dry_run", b()),
486 ],
487 &["node", "prop", "value"],
488 true,
489 ),
490 },
491 ToolSpec {
492 name: "sections_append",
493 description: "Append markdown at the end of a heading's section; `heading` is a heading block id or its text (scoped by `doc`/`path`).",
494 input_schema: schema(
495 &[
496 ("heading", s()),
497 ("markdown", s()),
498 ("doc", s()),
499 ("path", s()),
500 ("dry_run", b()),
501 ],
502 &["heading", "markdown"],
503 true,
504 ),
505 },
506 ToolSpec {
507 name: "docs_append",
508 description: "Append markdown at the end of a document as new top-level blocks (existing ids preserved).",
509 input_schema: schema(
510 &[("doc", s()), ("path", s()), ("text", s())],
511 &["text"],
512 true,
513 ),
514 },
515 ToolSpec {
516 name: "links_retarget",
517 description: "Rewrite one link destination everywhere it is linked (dry run by default).",
518 input_schema: schema(
519 &[
520 ("from_target", s()),
521 ("to_target", s()),
522 ("path_glob", s()),
523 ("dry_run", b()),
524 ],
525 &["from_target", "to_target"],
526 true,
527 ),
528 },
529 ToolSpec {
530 name: "links_stale",
531 description: "Dangling internal links (open edges to a `phantom:` target) plus external and total counts; `summary:true` returns counts only.",
532 input_schema: schema(
533 &[("path_glob", s()), ("limit", i()), ("summary", b())],
534 &[],
535 true,
536 ),
537 },
538 ToolSpec {
539 name: "links_repair",
540 description: "Bulk link repair: `repairs` ([{from,to}]) or one `from_target`/`to_target` pair, in one changeset (dry run by default).",
541 input_schema: schema(
542 &[
543 (
544 "repairs",
545 json!({ "type": "array", "items": { "type": "object", "properties": { "from": { "type": "string" }, "to": { "type": "string" } }, "required": ["from", "to"] } }),
546 ),
547 ("from_target", s()),
548 ("to_target", s()),
549 ("path_glob", s()),
550 ("dry_run", b()),
551 ],
552 &[],
553 true,
554 ),
555 },
556 ToolSpec {
557 name: "docs_create",
558 description: "Create a document at `path` from `markdown` with optional `frontmatter`; `dry_run` returns the would-be file under `diffs` without writing.",
559 input_schema: schema(
560 &[
561 ("path", s()),
562 ("markdown", s()),
563 ("frontmatter", obj()),
564 ("dry_run", b()),
565 ],
566 &["path", "markdown"],
567 true,
568 ),
569 },
570 ToolSpec {
571 name: "docs_move",
572 description: "Rename a document to `to_path`, identity preserved; `retarget_inbound` rewrites inbound links; `dry_run` returns the plan (`dangling`, `retargeted`, per-file `diffs`) without renaming.",
573 input_schema: schema(
574 &[
575 ("doc", s()),
576 ("to_path", s()),
577 ("retarget_inbound", b()),
578 ("dry_run", b()),
579 ],
580 &["doc", "to_path"],
581 true,
582 ),
583 },
584 ToolSpec {
585 name: "docs_delete",
586 description: "Delete a document: tombstone it and remove the file; `dry_run` returns the removal under `diffs` without deleting.",
587 input_schema: schema(&[("doc", s()), ("dry_run", b())], &["doc"], true),
588 },
589 ToolSpec {
590 name: "docs_set_meta",
591 description: "Set and/or unset frontmatter keys, re-ingesting the document; `dry_run` returns the rewritten file under `diffs` without writing.",
592 input_schema: schema(
593 &[
594 ("doc", s()),
595 ("set", obj()),
596 ("unset", strings()),
597 ("dry_run", b()),
598 ],
599 &["doc"],
600 true,
601 ),
602 },
603 ToolSpec {
604 name: "docs_plan_update",
605 description: "Plan a whole-document update without applying: the opset and a one-line-per-op plan.",
606 input_schema: schema(&[("doc", s()), ("content", s())], &["doc", "content"], true),
607 },
608 ToolSpec {
609 name: "docs_update",
610 description: "Whole-document update with identity preservation: plan then apply (`dry_run` returns the plan only).",
611 input_schema: schema(
612 &[
613 ("doc", s()),
614 ("content", s()),
615 ("reason", s()),
616 ("dry_run", b()),
617 ],
618 &["doc", "content"],
619 true,
620 ),
621 },
622 ToolSpec {
623 name: "observe",
624 description: "Record `content` as the authoritative bytes at `path` (an observed-origin commit; an echo when unchanged).",
625 input_schema: schema(
626 &[("path", s()), ("content", s())],
627 &["path", "content"],
628 true,
629 ),
630 },
631 ToolSpec {
632 name: "observe_many",
633 description: "Observe several files under one timestamp and one pool sweep.",
634 input_schema: schema(
635 &[(
636 "files",
637 json!({ "type": "array", "items": { "type": "object", "properties": { "path": { "type": "string" }, "content": { "type": "string" } }, "required": ["path", "content"] } }),
638 )],
639 &["files"],
640 true,
641 ),
642 },
643 ToolSpec {
644 name: "observe_delete",
645 description: "Record that `path` left the source: an observed, pooled tombstone.",
646 input_schema: schema(&[("path", s())], &["path"], true),
647 },
648 ToolSpec {
649 name: "history_node",
650 description: "A block's biography: the commits that touched it, newest first.",
651 input_schema: schema(&[("id", s()), ("limit", i())], &["id"], false),
652 },
653 ToolSpec {
654 name: "diff",
655 description: "Block-grain diff between two revisions of a document.",
656 input_schema: schema(
657 &[("doc", s()), ("from_rev", s()), ("to_rev", s())],
658 &["doc", "from_rev", "to_rev"],
659 true,
660 ),
661 },
662 ToolSpec {
663 name: "diff_unified",
664 description: "Unified diff (Myers, 3 lines of context) between two revisions (default: the previous and current).",
665 input_schema: schema(
666 &[("doc", s()), ("from_rev", s()), ("to_rev", s())],
667 &["doc"],
668 true,
669 ),
670 },
671 ToolSpec {
672 name: "docs_read_at",
673 description: "The whole document as of a past revision.",
674 input_schema: dp(&[("rev", s())], &["rev"]),
675 },
676 ToolSpec {
677 name: "docs_history",
678 description: "Version history of the docs matching `path_glob` or `doc`, grouped by document.",
679 input_schema: schema(
680 &[
681 ("path_glob", s()),
682 ("doc", s()),
683 ("include_deleted", b()),
684 ("limit", i()),
685 ],
686 &[],
687 true,
688 ),
689 },
690 ToolSpec {
691 name: "changes_since",
692 description: "The change feed: commit digests after `cursor` (a repo commit seq).",
693 input_schema: schema(
694 &[
695 ("cursor", i()),
696 (
697 "origin",
698 json!({ "type": "string", "enum": ["api", "observed", "import"] }),
699 ),
700 ("limit", i()),
701 ],
702 &[],
703 true,
704 ),
705 },
706 ToolSpec {
707 name: "repos_status",
708 description: "Repo counts, unconverged docs and on-disk drift.",
709 input_schema: schema(&[], &[], true),
710 },
711 ToolSpec {
712 name: "sync_status",
713 description: "Sync state: last commit seq, last checkpoint, convergence.",
714 input_schema: schema(&[], &[], true),
715 },
716 ToolSpec {
717 name: "repos",
718 description: "The repos in this workspace: `{ repos: [{ slug, hasSource }] }`.",
719 input_schema: schema(&[], &[], false),
720 },
721 ]
722}
723
724impl Surface {
725 #[must_use]
728 pub fn new(
729 store: Store,
730 default_repo: &str,
731 provider: Option<Box<dyn EmbeddingProvider>>,
732 ) -> Self {
733 let _ = store
744 .conn()
745 .execute_batch("PRAGMA cache_size = -16000; PRAGMA temp_store = MEMORY;");
746 Self {
747 store,
748 default_repo: default_repo.to_owned(),
749 provider,
750 on_mutation: None,
751 clock: Box::new(omgbase_sync::now_ts),
752 writes: WriteTarget::Derived,
753 config: Config::default(),
754 actor: ACTOR.to_owned(),
755 }
756 }
757
758 #[must_use]
761 pub fn with_actor(mut self, actor: &str) -> Self {
762 self.set_actor(actor);
763 self
764 }
765
766 pub fn set_actor(&mut self, actor: &str) {
768 self.actor = actor.to_owned();
769 }
770
771 #[must_use]
774 pub fn with_doc_store(mut self, doc_store: Box<dyn DocStore>) -> Self {
775 self.writes = WriteTarget::Fixed(doc_store);
776 self
777 }
778
779 #[must_use]
781 pub fn with_clock(mut self, clock: impl FnMut() -> String + 'static) -> Self {
782 self.clock = Box::new(clock);
783 self
784 }
785
786 #[must_use]
788 pub fn with_mutation_hook(mut self, hook: impl FnMut() + 'static) -> Self {
789 self.on_mutation = Some(Box::new(hook));
790 self
791 }
792
793 #[must_use]
795 pub fn with_config(mut self, config: Config) -> Self {
796 self.config = config;
797 self
798 }
799
800 #[must_use]
801 pub fn store(&self) -> &Store {
802 &self.store
803 }
804
805 pub fn store_mut(&mut self) -> &mut Store {
806 &mut self.store
807 }
808
809 #[must_use]
810 pub fn default_repo(&self) -> &str {
811 &self.default_repo
812 }
813
814 #[must_use]
816 pub fn tools(&self) -> Vec<ToolSpec> {
817 tools()
818 }
819
820 #[must_use]
824 pub fn is_write_tool(name: &str) -> bool {
825 matches!(
826 name,
827 "apply"
828 | "blocks_insert"
829 | "blocks_update"
830 | "blocks_move"
831 | "blocks_remove"
832 | "blocks_split"
833 | "blocks_merge"
834 | "tasks_complete"
835 | "node_set"
836 | "sections_append"
837 | "docs_append"
838 | "links_retarget"
839 | "links_repair"
840 | "docs_create"
841 | "docs_move"
842 | "docs_delete"
843 | "docs_set_meta"
844 | "docs_update"
845 | "observe"
846 | "observe_many"
847 | "observe_delete"
848 )
849 }
850
851 pub fn call(&mut self, name: &str, args: Json) -> ToolOutcome {
853 match self.call_result(name, &args) {
854 Ok(body) => ToolOutcome {
855 body,
856 is_error: false,
857 },
858 Err(e) => ToolOutcome {
859 body: e.to_json(),
860 is_error: true,
861 },
862 }
863 }
864
865 fn now(&mut self) -> String {
866 (self.clock)()
867 }
868
869 fn notify(&mut self) {
870 if let Some(hook) = &mut self.on_mutation {
871 hook();
872 }
873 }
874
875 fn repo_rows(&self) -> Result<Vec<RepoRow>> {
878 Ok(omgbase_sync::workspace::list_repos(&self.store)?)
879 }
880
881 fn scope(&self, args: &Json) -> Result<(String, Option<String>)> {
883 let rows = self.repo_rows()?;
884 match arg_str(args, "repo") {
885 None => {
886 let root = rows
887 .iter()
888 .find(|r| r.repo_id == self.default_repo)
889 .and_then(|r| r.root_path.clone());
890 Ok((self.default_repo.clone(), root))
891 }
892 Some(slug) => rows
893 .iter()
894 .find(|r| r.slug == slug)
895 .map(|r| (r.repo_id.clone(), r.root_path.clone()))
896 .ok_or_else(|| {
897 SurfaceError::with_data(
898 "repo_not_found",
899 format!("no repo '{slug}' in this workspace"),
900 json!({ "repo": slug }),
901 )
902 }),
903 }
904 }
905
906 fn require_root(&self, root: Option<&str>) -> Result<()> {
907 if matches!(self.writes, WriteTarget::Fixed(_)) || root.is_some() {
908 Ok(())
909 } else {
910 Err(SurfaceError::new(
911 "repo_not_found",
912 "repo has no filesystem source; mutation disabled",
913 ))
914 }
915 }
916
917 fn with_writes<T>(
919 &mut self,
920 root: Option<&str>,
921 f: impl FnOnce(&mut Store, &mut dyn DocStore) -> Result<T>,
922 ) -> Result<T> {
923 self.require_root(root)?;
924 match &mut self.writes {
925 WriteTarget::Fixed(ds) => f(&mut self.store, ds.as_mut()),
926 WriteTarget::Derived => {
927 let mut fs = FsDocStore::new(root.expect("checked by require_root"));
928 f(&mut self.store, &mut fs)
929 }
930 }
931 }
932
933 fn resolve_doc_id(
937 &self,
938 repo_id: &str,
939 doc: Option<&str>,
940 path: Option<&str>,
941 block: Option<&str>,
942 ) -> Result<String> {
943 let conn = self.store.conn();
944 let found = if let Some(d) = doc.filter(|d| !d.is_empty()) {
945 read::find_doc_by_ref(conn, repo_id, d)?.map(|i| i.doc_id)
946 } else if let Some(p) = path.filter(|p| !p.is_empty()) {
947 read::find_doc_by_path(conn, repo_id, p)?.map(|i| i.doc_id)
948 } else if let Some(b) = block.filter(|b| !b.is_empty()) {
949 conn.query_row(
950 "SELECT doc_id FROM blocks WHERE block_id = ?1",
951 params![b],
952 |r| r.get::<_, String>(0),
953 )
954 .optional()?
955 } else {
956 None
957 };
958 found.ok_or_else(|| {
959 let mut m = Map::new();
960 if let Some(d) = doc {
961 m.insert("doc".to_owned(), json!(d));
962 }
963 if let Some(p) = path {
964 m.insert("path".to_owned(), json!(p));
965 }
966 if let Some(b) = block {
967 m.insert("block".to_owned(), json!(b));
968 }
969 SurfaceError::with_data(
970 "doc_missing",
971 format!("no document for {}", Json::Object(m.clone())),
972 Json::Object(m),
973 )
974 })
975 }
976
977 fn resolve_doc_from_args(&self, repo_id: &str, args: &Json) -> Result<String> {
978 self.resolve_doc_id(repo_id, arg_str(args, "doc"), arg_str(args, "path"), None)
979 }
980
981 fn resolve_heading_id(
983 &self,
984 repo_id: &str,
985 heading: &str,
986 doc: Option<&str>,
987 path: Option<&str>,
988 ) -> Result<String> {
989 let conn = self.store.conn();
990 let as_block: Option<String> = conn
991 .query_row(
992 "SELECT block_id FROM blocks WHERE block_id = ?1 AND type = 'heading' AND deleted_commit IS NULL",
993 params![heading],
994 |r| r.get(0),
995 )
996 .optional()?;
997 if let Some(b) = as_block {
998 return Ok(b);
999 }
1000 let want_doc = if doc.is_some_and(|d| !d.is_empty()) || path.is_some_and(|p| !p.is_empty())
1001 {
1002 Some(self.resolve_doc_id(repo_id, doc, path, None)?)
1003 } else {
1004 None
1005 };
1006 let needle = if heading.trim_start_matches([' ', '\t']).starts_with('#') {
1007 normalize_visible_text(heading, BlockKind::Heading, 0)
1008 } else {
1009 normalize_text(heading)
1010 };
1011 let rows: Vec<(String, String, String)> = match &want_doc {
1012 Some(d) => {
1013 let mut stmt = conn.prepare(
1014 "SELECT block_id, doc_id, text FROM blocks WHERE repo_id = ?1 AND type = 'heading' AND doc_id = ?2 AND deleted_commit IS NULL",
1015 )?;
1016 let it = stmt.query_map(params![repo_id, d], |r| {
1017 Ok((r.get(0)?, r.get(1)?, r.get(2)?))
1018 })?;
1019 it.collect::<std::result::Result<_, _>>()?
1020 }
1021 None => {
1022 let mut stmt = conn.prepare(
1023 "SELECT block_id, doc_id, text FROM blocks WHERE repo_id = ?1 AND type = 'heading' AND deleted_commit IS NULL",
1024 )?;
1025 let it =
1026 stmt.query_map(params![repo_id], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)))?;
1027 it.collect::<std::result::Result<_, _>>()?
1028 }
1029 };
1030 let matches: Vec<&(String, String, String)> = rows
1031 .iter()
1032 .filter(|(_, _, text)| normalize_text(text) == needle)
1033 .collect();
1034 match matches.len() {
1035 0 => {
1036 let mut data = Map::new();
1037 data.insert("heading".to_owned(), json!(heading));
1038 if let Some(d) = want_doc {
1039 data.insert("doc".to_owned(), json!(d));
1040 }
1041 Err(SurfaceError::with_data(
1042 "parent_missing",
1043 format!("no heading matching {}", Json::String(heading.to_owned())),
1044 Json::Object(data),
1045 ))
1046 }
1047 1 => Ok(matches[0].0.clone()),
1048 n => Err(SurfaceError::with_data(
1049 "ambiguous_heading",
1050 format!(
1051 "heading {} matches {n} headings; pass its block id or a doc/path scope",
1052 Json::String(heading.to_owned())
1053 ),
1054 json!({
1055 "heading": heading,
1056 "candidates": matches.iter().map(|(b, d, _)| json!({ "block": b, "doc": d })).collect::<Vec<_>>(),
1057 }),
1058 )),
1059 }
1060 }
1061
1062 fn resolve_block_ref(&self, repo_id: &str, r: &str) -> Result<String> {
1064 match read::resolve_ref(self.store.conn(), repo_id, r)? {
1065 Some(ResolvedRef::Block { block_id, .. }) => Ok(block_id),
1066 _ => Err(SurfaceError::with_data(
1067 "block_missing",
1068 format!("not a block: {r}"),
1069 json!({ "ref": r }),
1070 )),
1071 }
1072 }
1073
1074 fn resolve_parent_ref(&self, repo_id: &str, r: &str) -> Result<(Parent, String)> {
1076 match read::resolve_ref(self.store.conn(), repo_id, r)? {
1077 Some(ResolvedRef::Block { doc_id, block_id }) => Ok((Parent::Block(block_id), doc_id)),
1078 Some(ResolvedRef::Document { doc_id }) => Ok((Parent::Doc, doc_id)),
1079 None => Err(SurfaceError::with_data(
1080 "block_missing",
1081 format!("not a block or document: {r}"),
1082 json!({ "ref": r }),
1083 )),
1084 }
1085 }
1086
1087 fn doc_id_of_block(&self, block_id: &str) -> Result<Option<String>> {
1088 Ok(self
1089 .store
1090 .conn()
1091 .query_row(
1092 "SELECT doc_id FROM blocks WHERE block_id = ?1 AND deleted_commit IS NULL",
1093 params![block_id],
1094 |r| r.get(0),
1095 )
1096 .optional()?)
1097 }
1098
1099 fn pin_hash(&self, block_id: &str) -> Result<Option<String>> {
1101 let h: Option<Vec<u8>> = self
1102 .store
1103 .conn()
1104 .query_row(
1105 "SELECT raw_hash FROM blocks WHERE block_id = ?1 AND deleted_commit IS NULL",
1106 params![block_id],
1107 |r| r.get(0),
1108 )
1109 .optional()?;
1110 Ok(h.map(|h| omgbase_format::hash::hex(&h)))
1111 }
1112
1113 fn resolve_at(&self, repo_id: &str, at: Option<&Json>) -> Result<At> {
1115 match at {
1116 None | Some(Json::Null) => Ok(At::End),
1117 Some(Json::String(s)) if s == "start" => Ok(At::Start),
1118 Some(Json::String(s)) if s == "end" => Ok(At::End),
1119 Some(Json::Object(o)) => {
1120 if let Some(b) = o.get("before").and_then(Json::as_str) {
1121 return Ok(At::Before(self.resolve_block_ref(repo_id, b)?));
1122 }
1123 if let Some(a) = o.get("after").and_then(Json::as_str) {
1124 return Ok(At::After(self.resolve_block_ref(repo_id, a)?));
1125 }
1126 Err(bad_args(
1127 "`at` must be \"start\", \"end\", {before} or {after}",
1128 ))
1129 }
1130 Some(_) => Err(bad_args(
1131 "`at` must be \"start\", \"end\", {before} or {after}",
1132 )),
1133 }
1134 }
1135
1136 fn apply_ops(
1139 &mut self,
1140 repo_id: &str,
1141 root: Option<&str>,
1142 ops: Vec<Op>,
1143 reason: &str,
1144 dry_run: bool,
1145 ) -> Result<ApplyResult> {
1146 let ts = self.now();
1147 let req = ApplyRequest {
1148 repo_id: repo_id.to_owned(),
1149 ops,
1150 origin: ApplyOrigin::new(&self.actor, Some(reason)),
1151 dry_run,
1152 set_frontmatter: Vec::new(),
1153 };
1154 let res = self.with_writes(root, |store, ds| Ok(store.apply(&req, ds, &ts)?))?;
1155 if !dry_run {
1156 self.notify();
1157 }
1158 Ok(res)
1159 }
1160
1161 fn doc_ctx(&mut self, repo_id: &str) -> DocOpContext {
1162 DocOpContext {
1163 repo_id: repo_id.to_owned(),
1164 actor: Some(self.actor.clone()),
1165 ts: self.now(),
1166 }
1167 }
1168
1169 #[allow(clippy::too_many_lines)]
1173 pub fn call_result(&mut self, name: &str, args: &Json) -> Result<Json> {
1174 match name {
1175 "docs_outline" => {
1176 let (repo, _) = self.scope(args)?;
1177 let doc_id = self.resolve_doc_from_args(&repo, args)?;
1178 let skeleton = match arg_str(args, "resolution") {
1179 None | Some("outline") => false,
1180 Some("skeleton") => true,
1181 Some(_) => return Err(bad_args("`resolution` must be skeleton or outline")),
1182 };
1183 read::docs_outline(
1184 &self.store,
1185 &doc_id,
1186 skeleton,
1187 arg_i64(args, "depth")?,
1188 arg_usize(args, "budget_tokens")?,
1189 )
1190 }
1191 "docs_read" => {
1192 let (repo, _) = self.scope(args)?;
1193 let doc_id = self.resolve_doc_from_args(&repo, args)?;
1194 read::docs_read(
1195 &self.store,
1196 &doc_id,
1197 arg_bool(args, "include_ids")?.unwrap_or(false),
1198 )?
1199 .ok_or_else(|| SurfaceError::new("doc_missing", format!("no document for {args}")))
1200 }
1201 "docs_get_many" => {
1202 let (repo, _) = self.scope(args)?;
1203 let refs = arg_strings(args, "docs")?;
1204 read::docs_read_many(
1205 &self.store,
1206 &repo,
1207 &refs,
1208 arg_bool(args, "include_ids")?.unwrap_or(false),
1209 arg_usize(args, "budget_tokens")?,
1210 )
1211 }
1212 "nodes_get" => {
1213 let (repo, _) = self.scope(args)?;
1214 let id = arg_string(args, "id")?;
1215 let doc_id = self.resolve_doc_id(
1216 &repo,
1217 arg_str(args, "doc"),
1218 arg_str(args, "path"),
1219 Some(&id),
1220 )?;
1221 read::nodes_get(
1222 &self.store,
1223 &doc_id,
1224 &id,
1225 resolution_arg(args, Resolution::Full)?,
1226 )?
1227 .ok_or_else(|| SurfaceError::new("block_missing", format!("no block {id}")))
1228 }
1229 "nodes_get_many" => {
1230 let (repo, _) = self.scope(args)?;
1231 let ids = arg_strings(args, "ids")?;
1232 let scoped = arg_str(args, "doc").is_some_and(|d| !d.is_empty())
1233 || arg_str(args, "path").is_some_and(|p| !p.is_empty());
1234 let doc_id = if scoped {
1235 Some(self.resolve_doc_from_args(&repo, args)?)
1236 } else {
1237 None
1238 };
1239 read::nodes_get_many(
1240 &self.store,
1241 doc_id.as_deref(),
1242 &ids,
1243 resolution_arg(args, Resolution::Text)?,
1244 arg_usize(args, "budget_tokens")?,
1245 )
1246 }
1247 "read_ref" => {
1248 let (repo, _) = self.scope(args)?;
1249 let r = arg_string(args, "ref")?;
1250 let resolved =
1251 read::resolve_ref(self.store.conn(), &repo, &r)?.ok_or_else(|| {
1252 SurfaceError::new(
1253 "doc_missing",
1254 format!("no document or block for {}", Json::String(r.clone())),
1255 )
1256 })?;
1257 match resolved {
1258 ResolvedRef::Document { doc_id } => {
1259 let res =
1260 read::docs_read(&self.store, &doc_id, false)?.ok_or_else(|| {
1261 SurfaceError::new(
1262 "doc_missing",
1263 format!("no document for {}", Json::String(r.clone())),
1264 )
1265 })?;
1266 let mut m = Map::new();
1267 m.insert("kind".to_owned(), json!("document"));
1268 Ok(merge(m, res))
1269 }
1270 ResolvedRef::Block { doc_id, block_id } => {
1271 let node = read::nodes_get(
1272 &self.store,
1273 &doc_id,
1274 &block_id,
1275 resolution_arg(args, Resolution::Raw)?,
1276 )?
1277 .ok_or_else(|| {
1278 SurfaceError::new("block_missing", format!("no block {block_id}"))
1279 })?;
1280 let mut m = Map::new();
1281 m.insert("kind".to_owned(), json!("block"));
1282 Ok(merge(m, node))
1283 }
1284 }
1285 }
1286 "docs_tree" => {
1287 let (repo, _) = self.scope(args)?;
1288 read::docs_tree(
1289 &self.store,
1290 &repo,
1291 arg_str(args, "path"),
1292 arg_i64(args, "depth")?,
1293 arg_i64(args, "limit")?,
1294 arg_str(args, "cursor"),
1295 arg_usize(args, "budget_tokens")?,
1296 )
1297 }
1298 "docs_list" => {
1299 let (repo, _) = self.scope(args)?;
1300 read::docs_list(
1301 &self.store,
1302 &repo,
1303 arg_str(args, "path_glob"),
1304 arg_i64(args, "limit")?,
1305 arg_str(args, "cursor"),
1306 arg_usize(args, "budget_tokens")?,
1307 )
1308 }
1309 "query_syntax" => Ok(json!({ "syntax": QUERY_SYNTAX })),
1310 "query" => {
1311 let (repo, _) = self.scope(args)?;
1312 let source = arg_string(args, "query")?;
1313 if self.provider.is_none()
1316 && !crate::query::collect_semantic_phrases(&source).is_empty()
1317 {
1318 return Err(SurfaceError::new(
1319 "semantic_unavailable",
1320 "no embedding provider configured for this server",
1321 ));
1322 }
1323 let opts = QueryOptions {
1324 limit: arg_usize(args, "limit")?,
1325 cursor: arg_str(args, "cursor"),
1326 provider: self.provider.as_deref(),
1327 in_memory: false,
1328 };
1329 Ok(query(&self.store, &repo, &source, opts)?.to_json())
1330 }
1331 "graph" => {
1332 let (repo, _) = self.scope(args)?;
1333 let g = GraphArgs {
1334 roots: arg_strings(args, "roots")?,
1335 degrees: arg_i64(args, "degrees")?,
1336 direction: arg_str(args, "direction").map(str::to_owned),
1337 predicate: arg_str(args, "predicate").map(str::to_owned),
1338 select: if args.get("select").is_some() {
1339 arg_strings(args, "select")?
1340 } else {
1341 Vec::new()
1342 },
1343 max_documents: arg_i64(args, "max_documents")?,
1344 };
1345 if self.provider.is_none()
1346 && g.select.iter().any(|s| {
1347 !crate::query::collect_semantic_phrases(&format!("from docs select x: {s}"))
1348 .is_empty()
1349 })
1350 {
1351 return Err(SurfaceError::new(
1352 "semantic_unavailable",
1353 "no embedding provider configured for this server",
1354 ));
1355 }
1356 graph_neighborhood(&self.store, &repo, &g, self.provider.as_deref())
1357 }
1358 "text_search" => {
1359 let (repo, _) = self.scope(args)?;
1360 let q = arg_string(args, "q")?;
1361 let res =
1362 self.store
1363 .text_search(&repo, &q, arg_usize(args, "limit")?.unwrap_or(50))?;
1364 Ok(json!({
1365 "hits": res.hits.iter().map(|h| json!({
1366 "blockId": h.block_id, "docId": h.doc_id, "path": h.path, "type": h.block_type, "text": h.text, "score": h.score,
1367 })).collect::<Vec<_>>(),
1368 "truncated": res.truncated,
1369 }))
1370 }
1371 "resolve" => {
1372 let (repo, _) = self.scope(args)?;
1373 let q = arg_string(args, "query")?;
1374 let vector = match &self.provider {
1375 Some(p) => Some(QueryVector {
1376 model: p.model().to_owned(),
1377 vec: p
1378 .embed_query(&q)
1379 .map_err(|e| SurfaceError::new(e.code(), e.to_string()))?,
1380 }),
1381 None => None,
1382 };
1383 let hits = self
1384 .store
1385 .resolve(&repo, &q, vector, arg_usize(args, "limit")?)?;
1386 Ok(Json::Array(hits.iter().map(|h| json!({
1387 "id": h.id, "locator": h.locator, "preview": h.preview, "evidence": evidence_json(&h.evidence),
1388 })).collect()))
1389 }
1390 "apply" => {
1391 let (repo, root) = self.scope(args)?;
1392 self.require_root(root.as_deref())?;
1393 let ops = args
1394 .get("ops")
1395 .and_then(Json::as_array)
1396 .ok_or_else(|| bad_args("`ops` must be an array"))?;
1397 let ops: Vec<Op> = ops
1398 .iter()
1399 .map(|o| Op::from_json(o).map_err(bad_args))
1400 .collect::<Result<_>>()?;
1401 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1402 let reason = arg_str(args, "reason").map(str::to_owned);
1403 let ts = self.now();
1404 let req = ApplyRequest {
1405 repo_id: repo.clone(),
1406 ops,
1407 origin: ApplyOrigin::new(&self.actor, reason.as_deref()),
1408 dry_run: dry,
1409 set_frontmatter: Vec::new(),
1410 };
1411 let res =
1412 self.with_writes(root.as_deref(), |store, ds| Ok(store.apply(&req, ds, &ts)?))?;
1413 if !dry {
1414 self.notify();
1415 }
1416 Ok(apply_json(&res))
1417 }
1418 "blocks_insert" => {
1419 let (repo, root) = self.scope(args)?;
1420 self.require_root(root.as_deref())?;
1421 let (parent, doc_id) = self.resolve_parent_ref(&repo, &arg_string(args, "to")?)?;
1422 let at = self.resolve_at(&repo, args.get("at"))?;
1423 let doc = if parent == Parent::Doc {
1424 Some(doc_id)
1425 } else {
1426 None
1427 };
1428 let ops = vec![Op::Insert {
1429 doc,
1430 to: To { parent, at },
1431 markdown: arg_string(args, "markdown")?,
1432 expect: arg_parent_expect(args)?,
1433 }];
1434 Ok(apply_json(&self.apply_ops(
1435 &repo,
1436 root.as_deref(),
1437 ops,
1438 "blocks_insert",
1439 arg_bool(args, "dry_run")?.unwrap_or(false),
1440 )?))
1441 }
1442 "blocks_update" => {
1443 let (repo, root) = self.scope(args)?;
1444 self.require_root(root.as_deref())?;
1445 let block = self.resolve_block_ref(&repo, &arg_string(args, "block")?)?;
1446 let expect = match arg_expect(args)? {
1447 Some(e) => Some(e),
1448 None => self.pin_hash(&block)?.map(Expect::content),
1449 };
1450 let checked = arg_bool(args, "checked")?;
1451 let extra = arg_object(args, "attrs")?;
1452 let attrs = if checked.is_some() || extra.is_some() {
1453 let mut a = Map::new();
1454 if let Some(c) = checked {
1455 a.insert("checked".to_owned(), json!(c));
1456 }
1457 if let Some(x) = extra {
1458 for (k, v) in x {
1459 a.insert(k.clone(), v.clone());
1460 }
1461 }
1462 Some(a)
1463 } else {
1464 None
1465 };
1466 let ops = vec![Op::Update {
1467 block: block.clone(),
1468 markdown: arg_str(args, "markdown").map(str::to_owned),
1469 attrs,
1470 expect,
1471 trivia: None,
1472 child_ids: None,
1473 }];
1474 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1475 let res = self.apply_ops(&repo, root.as_deref(), ops, "blocks_update", dry)?;
1476 let ids: Vec<String> = res
1477 .results
1478 .first()
1479 .map_or_else(|| vec![block.clone()], |r| r.ids.clone());
1480 let mut m = Map::new();
1481 m.insert(
1482 "id".to_owned(),
1483 json!(ids.first().cloned().unwrap_or(block)),
1484 );
1485 m.insert("ids".to_owned(), json!(ids));
1486 Ok(merge(m, apply_json(&res)))
1487 }
1488 "blocks_move" => {
1489 let (repo, root) = self.scope(args)?;
1490 self.require_root(root.as_deref())?;
1491 let blocks: Vec<String> = arg_strings(args, "blocks")?
1492 .iter()
1493 .map(|b| self.resolve_block_ref(&repo, b))
1494 .collect::<Result<_>>()?;
1495 let to_ref = arg_string(args, "to")?;
1496 let (parent, doc_id) = self.resolve_parent_ref(&repo, &to_ref)?;
1497 if parent == Parent::Doc {
1498 if let Some(first) = blocks.first() {
1499 if self.doc_id_of_block(first)? != Some(doc_id) {
1500 return Err(SurfaceError::with_data(
1501 "target_missing",
1502 format!(
1503 "blocks_move cannot target another document's root ({to_ref}); anchor on a block in that document with at.before/at.after"
1504 ),
1505 json!({ "to": to_ref }),
1506 ));
1507 }
1508 }
1509 }
1510 let at = self.resolve_at(&repo, args.get("at"))?;
1511 let ops = vec![Op::Move {
1512 blocks,
1513 to: To { parent, at },
1514 expect: arg_parent_expect(args)?,
1515 }];
1516 Ok(apply_json(&self.apply_ops(
1517 &repo,
1518 root.as_deref(),
1519 ops,
1520 "blocks_move",
1521 arg_bool(args, "dry_run")?.unwrap_or(false),
1522 )?))
1523 }
1524 "blocks_remove" => {
1525 let (repo, root) = self.scope(args)?;
1526 self.require_root(root.as_deref())?;
1527 let blocks: Vec<String> = arg_strings(args, "blocks")?
1528 .iter()
1529 .map(|b| self.resolve_block_ref(&repo, b))
1530 .collect::<Result<_>>()?;
1531 let ops = vec![Op::Remove {
1532 blocks,
1533 expect: None,
1534 }];
1535 Ok(apply_json(&self.apply_ops(
1536 &repo,
1537 root.as_deref(),
1538 ops,
1539 "blocks_remove",
1540 arg_bool(args, "dry_run")?.unwrap_or(false),
1541 )?))
1542 }
1543 "blocks_split" => {
1544 let (repo, root) = self.scope(args)?;
1545 self.require_root(root.as_deref())?;
1546 let block = self.resolve_block_ref(&repo, &arg_string(args, "block")?)?;
1547 let at: Vec<usize> = args
1548 .get("at")
1549 .and_then(Json::as_array)
1550 .ok_or_else(|| bad_args("`at` must be an array of byte offsets"))?
1551 .iter()
1552 .map(|v| {
1553 v.as_u64()
1554 .map(|n| usize::try_from(n).unwrap_or(usize::MAX))
1555 .ok_or_else(|| bad_args("`at` must be an array of byte offsets"))
1556 })
1557 .collect::<Result<_>>()?;
1558 let expect = Some(Expect::content(self.pin_hash(&block)?.unwrap_or_default()));
1560 let ops = vec![Op::Split { block, at, expect }];
1561 Ok(apply_json(&self.apply_ops(
1562 &repo,
1563 root.as_deref(),
1564 ops,
1565 "blocks_split",
1566 arg_bool(args, "dry_run")?.unwrap_or(false),
1567 )?))
1568 }
1569 "blocks_merge" => {
1570 let (repo, root) = self.scope(args)?;
1571 self.require_root(root.as_deref())?;
1572 let blocks: Vec<String> = arg_strings(args, "blocks")?
1573 .iter()
1574 .map(|b| self.resolve_block_ref(&repo, b))
1575 .collect::<Result<_>>()?;
1576 let ops = vec![Op::Merge {
1577 blocks,
1578 separator: arg_str(args, "separator").map(str::to_owned),
1579 expect: None,
1580 }];
1581 Ok(apply_json(&self.apply_ops(
1582 &repo,
1583 root.as_deref(),
1584 ops,
1585 "blocks_merge",
1586 arg_bool(args, "dry_run")?.unwrap_or(false),
1587 )?))
1588 }
1589 "tasks_complete" => {
1590 let (repo, root) = self.scope(args)?;
1591 self.require_root(root.as_deref())?;
1592 let blocks: Vec<String> = arg_strings(args, "blocks")?
1593 .iter()
1594 .map(|b| self.resolve_block_ref(&repo, b))
1595 .collect::<Result<_>>()?;
1596 let ops = if arg_bool(args, "checked")?.unwrap_or(true) {
1597 self.store.tasks_complete(&blocks)?
1598 } else {
1599 let mut ops = Vec::with_capacity(blocks.len());
1600 for b in &blocks {
1601 let mut a = Map::new();
1602 a.insert("checked".to_owned(), json!(false));
1603 ops.push(Op::Update {
1604 block: b.clone(),
1605 markdown: None,
1606 attrs: Some(a),
1607 expect: self.pin_hash(b)?.map(Expect::content),
1608 trivia: None,
1609 child_ids: None,
1610 });
1611 }
1612 ops
1613 };
1614 Ok(apply_json(&self.apply_ops(
1615 &repo,
1616 root.as_deref(),
1617 ops,
1618 "tasks_complete",
1619 arg_bool(args, "dry_run")?.unwrap_or(false),
1620 )?))
1621 }
1622 "node_set" => {
1623 let (repo, root) = self.scope(args)?;
1624 self.require_root(root.as_deref())?;
1625 let ops = self.store.node_set(
1626 &arg_string(args, "node")?,
1627 &arg_string(args, "prop")?,
1628 &arg_string(args, "value")?,
1629 )?;
1630 Ok(apply_json(&self.apply_ops(
1631 &repo,
1632 root.as_deref(),
1633 ops,
1634 "node_set",
1635 arg_bool(args, "dry_run")?.unwrap_or(false),
1636 )?))
1637 }
1638 "sections_append" => {
1639 let (repo, root) = self.scope(args)?;
1640 self.require_root(root.as_deref())?;
1641 let heading = self.resolve_heading_id(
1642 &repo,
1643 &arg_string(args, "heading")?,
1644 arg_str(args, "doc"),
1645 arg_str(args, "path"),
1646 )?;
1647 let ops = Store::sections_append(&heading, &arg_string(args, "markdown")?);
1648 Ok(apply_json(&self.apply_ops(
1649 &repo,
1650 root.as_deref(),
1651 ops,
1652 "sections_append",
1653 arg_bool(args, "dry_run")?.unwrap_or(false),
1654 )?))
1655 }
1656 "docs_append" => {
1657 let (repo, root) = self.scope(args)?;
1658 self.require_root(root.as_deref())?;
1659 let doc_id = self.resolve_doc_from_args(&repo, args)?;
1660 let ops = Store::docs_append(&doc_id, &arg_string(args, "text")?);
1661 Ok(apply_json(&self.apply_ops(
1662 &repo,
1663 root.as_deref(),
1664 ops,
1665 "docs_append",
1666 false,
1667 )?))
1668 }
1669 "links_retarget" | "links_repair" => {
1670 let (repo, root) = self.scope(args)?;
1671 self.require_root(root.as_deref())?;
1672 let repairs: Vec<omgbase_store::LinkRepair> = if name == "links_retarget" {
1673 vec![omgbase_store::LinkRepair {
1674 from: arg_string(args, "from_target")?,
1675 to: arg_string(args, "to_target")?,
1676 }]
1677 } else if let Some(list) = args.get("repairs").and_then(Json::as_array) {
1678 list.iter()
1679 .map(|r| {
1680 Ok(omgbase_store::LinkRepair {
1681 from: arg_string(r, "from")?,
1682 to: arg_string(r, "to")?,
1683 })
1684 })
1685 .collect::<Result<_>>()?
1686 } else if let (Some(f), Some(t)) =
1687 (arg_str(args, "from_target"), arg_str(args, "to_target"))
1688 {
1689 vec![omgbase_store::LinkRepair {
1690 from: f.to_owned(),
1691 to: t.to_owned(),
1692 }]
1693 } else {
1694 Vec::new()
1695 };
1696 if repairs.is_empty() {
1697 return Err(SurfaceError::new(
1698 "target_missing",
1699 "links_repair requires `repairs` (array of {from,to}) or a `from_target`+`to_target` pair",
1700 ));
1701 }
1702 let plan = self.store.links_repair(
1703 &repo,
1704 &repairs,
1705 arg_str(args, "path_glob").filter(|g| !g.is_empty()),
1706 )?;
1707 let dry = arg_bool(args, "dry_run")?.unwrap_or(true);
1708 let res = self.apply_ops(&repo, root.as_deref(), plan.ops.clone(), name, dry)?;
1709 let mut m = Map::new();
1710 m.insert(
1711 "hits".to_owned(),
1712 Json::Array(plan.hits.iter().map(|h| json!({ "block": h.block, "path": h.path, "oldRaw": h.old_raw, "newRaw": h.new_raw })).collect()),
1713 );
1714 m.insert(
1715 "pairs".to_owned(),
1716 Json::Array(
1717 plan.pairs
1718 .iter()
1719 .map(|p| json!({ "from": p.from, "to": p.to, "hits": p.hits }))
1720 .collect(),
1721 ),
1722 );
1723 m.insert("applied".to_owned(), json!(!dry));
1724 Ok(merge(m, apply_json(&res)))
1725 }
1726 "links_stale" => {
1727 let (repo, _) = self.scope(args)?;
1728 let glob = arg_str(args, "path_glob").filter(|g| !g.is_empty());
1729 if arg_bool(args, "summary")?.unwrap_or(false) {
1730 return links::links_stale_summary(&self.store, &repo, glob);
1731 }
1732 links::links_stale(&self.store, &repo, glob, arg_i64(args, "limit")?)
1733 }
1734 "docs_create" => {
1735 let (repo, root) = self.scope(args)?;
1736 let ctx = self.doc_ctx(&repo);
1737 let path = arg_string(args, "path")?;
1738 let markdown = arg_string(args, "markdown")?;
1739 let fm = arg_object(args, "frontmatter")?.cloned();
1740 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1741 let res = self.with_writes(root.as_deref(), |store, ds| {
1742 Ok(if dry {
1743 store
1744 .dry_run()
1745 .docs_create(&ctx, ds, &path, &markdown, fm.as_ref())?
1746 } else {
1747 store.docs_create(&ctx, ds, &path, &markdown, fm.as_ref())?
1748 })
1749 })?;
1750 if !dry {
1751 self.notify();
1752 }
1753 Ok(doc_op_json(&res))
1754 }
1755 "docs_move" => {
1756 let (repo, root) = self.scope(args)?;
1757 let ctx = self.doc_ctx(&repo);
1758 let doc = arg_string(args, "doc")?;
1759 let to = arg_string(args, "to_path")?;
1760 let retarget = arg_bool(args, "retarget_inbound")?.unwrap_or(false);
1761 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1762 let res = self.with_writes(root.as_deref(), |store, ds| {
1763 Ok(if dry {
1764 store.dry_run().docs_move(&ctx, ds, &doc, &to, retarget)?
1765 } else {
1766 store.docs_move(&ctx, ds, &doc, &to, retarget)?
1767 })
1768 })?;
1769 if !dry {
1770 self.notify();
1771 }
1772 let mut m = Map::new();
1775 m.insert("docId".to_owned(), json!(res.doc_id));
1776 m.insert("path".to_owned(), json!(res.path));
1777 m.insert("committed".to_owned(), json!(res.committed));
1778 if let Some(diffs) = &res.diffs {
1779 m.insert("diffs".to_owned(), diffs_json(diffs));
1780 }
1781 m.insert(
1782 "dangling".to_owned(),
1783 Json::Array(
1784 res.dangling
1785 .iter()
1786 .map(omgbase_store::InboundLink::to_json)
1787 .collect(),
1788 ),
1789 );
1790 m.insert(
1791 "retargeted".to_owned(),
1792 json!(
1793 res.retargeted
1794 .as_ref()
1795 .map(|r| json!({ "blocks": r.blocks, "docs": r.docs }))
1796 ),
1797 );
1798 Ok(Json::Object(m))
1799 }
1800 "docs_delete" => {
1801 let (repo, root) = self.scope(args)?;
1802 let ctx = self.doc_ctx(&repo);
1803 let doc = arg_string(args, "doc")?;
1804 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1805 let res = self.with_writes(root.as_deref(), |store, ds| {
1806 Ok(if dry {
1807 store.dry_run().docs_delete(&ctx, ds, &doc)?
1808 } else {
1809 store.docs_delete(&ctx, ds, &doc)?
1810 })
1811 })?;
1812 if !dry {
1813 self.notify();
1814 }
1815 Ok(doc_op_json(&res))
1816 }
1817 "docs_set_meta" => {
1818 let (repo, root) = self.scope(args)?;
1819 let ctx = self.doc_ctx(&repo);
1820 let doc = arg_string(args, "doc")?;
1821 let set = arg_object(args, "set")?.cloned();
1822 let unset = if args.get("unset").is_some() {
1823 arg_strings(args, "unset")?
1824 } else {
1825 Vec::new()
1826 };
1827 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1828 let res = self.with_writes(root.as_deref(), |store, ds| {
1829 Ok(if dry {
1830 store
1831 .dry_run()
1832 .docs_set_meta(&ctx, ds, &doc, set.as_ref(), &unset)?
1833 } else {
1834 store.docs_set_meta(&ctx, ds, &doc, set.as_ref(), &unset)?
1835 })
1836 })?;
1837 if !dry {
1838 self.notify();
1839 }
1840 Ok(doc_op_json(&res))
1841 }
1842 "docs_plan_update" => {
1843 let (repo, root) = self.scope(args)?;
1844 self.require_root(root.as_deref())?;
1845 let doc = arg_string(args, "doc")?;
1846 let content = arg_string(args, "content")?;
1847 let opset = self
1848 .store
1849 .plan_update(&repo, &doc, &content, &self.config.clone())?;
1850 Ok(json!({ "opset": opset_json(&opset), "plan": render_opset_plan(&opset) }))
1851 }
1852 "docs_update" => {
1853 let (repo, root) = self.scope(args)?;
1854 let doc = arg_string(args, "doc")?;
1855 let content = arg_string(args, "content")?;
1856 let dry = arg_bool(args, "dry_run")?.unwrap_or(false);
1857 let reason = arg_str(args, "reason").map(str::to_owned);
1858 let origin = ApplyOrigin::new(&self.actor, reason.as_deref());
1859 let config = self.config.clone();
1860 let ts = self.now();
1861 let (opset, result) = self.with_writes(root.as_deref(), |store, ds| {
1862 Ok(store.docs_update(&repo, &doc, &content, &config, &origin, dry, ds, &ts)?)
1863 })?;
1864 if !dry {
1865 self.notify();
1866 }
1867 Ok(json!({
1868 "opset": opset_json(&opset),
1869 "plan": render_opset_plan(&opset),
1870 "result": result.map(|r| apply_json(&r)),
1871 }))
1872 }
1873 "observe" => {
1874 let (repo, _) = self.scope(args)?;
1875 let path = arg_string(args, "path")?;
1876 let content = arg_string(args, "content")?;
1877 let ts = self.now();
1878 let config = self.config.clone();
1879 let out = self
1880 .store
1881 .observe_one(&repo, &path, &content, &ts, &config)?;
1882 self.store.sweep_pool(&ts)?;
1883 self.notify();
1884 Ok(observe_json(&out))
1885 }
1886 "observe_many" => {
1887 let (repo, _) = self.scope(args)?;
1888 let files = args
1889 .get("files")
1890 .and_then(Json::as_array)
1891 .ok_or_else(|| bad_args("`files` must be an array of {path, content}"))?;
1892 let items: Vec<omgbase_store::BatchItem> = files
1893 .iter()
1894 .map(|f| {
1895 Ok(omgbase_store::BatchItem::observed(
1896 &arg_string(f, "path")?,
1897 &arg_string(f, "content")?,
1898 ))
1899 })
1900 .collect::<Result<_>>()?;
1901 let ts = self.now();
1902 let config = self.config.clone();
1903 let outcomes = self.store.observe_batch(&repo, &items, &ts, &config)?;
1904 self.store.sweep_pool(&ts)?;
1905 self.notify();
1906 let mut out = Vec::with_capacity(outcomes.len());
1907 for o in &outcomes {
1908 match o.as_observed() {
1909 Some(obs) => out.push(observe_json(obs)),
1910 None => {
1911 return Err(SurfaceError::other(format!(
1912 "observe_many: unexpected outcome for {}",
1913 o.path()
1914 )));
1915 }
1916 }
1917 }
1918 Ok(Json::Array(out))
1919 }
1920 "observe_delete" => {
1921 let (repo, _) = self.scope(args)?;
1922 let path = arg_string(args, "path")?;
1923 let ts = self.now();
1924 let out = self.store.observe_delete(&repo, &path, &ts)?;
1925 self.notify();
1926 Ok(json!({ "docId": out.doc_id, "path": out.path, "deleted": out.deleted() }))
1927 }
1928 "history_node" => history::history_node(
1929 &self.store,
1930 &arg_string(args, "id")?,
1931 arg_i64(args, "limit")?,
1932 ),
1933 "diff" => {
1934 let (repo, _) = self.scope(args)?;
1935 let doc_id = self.resolve_doc_id(&repo, arg_str(args, "doc"), None, None)?;
1936 history::diff_blocks(
1937 &self.store,
1938 &doc_id,
1939 &arg_string(args, "from_rev")?,
1940 &arg_string(args, "to_rev")?,
1941 )
1942 }
1943 "diff_unified" => {
1944 let (repo, _) = self.scope(args)?;
1945 let doc_ref = arg_string(args, "doc")?;
1946 let doc_id = self.resolve_doc_id(&repo, Some(&doc_ref), None, None)?;
1947 let revs = history::recent_revs(self.store.conn(), &doc_id)?;
1948 let to_rev = arg_str(args, "to_rev")
1949 .map(str::to_owned)
1950 .or_else(|| revs.first().cloned());
1951 let from_rev = arg_str(args, "from_rev")
1952 .map(str::to_owned)
1953 .or_else(|| revs.get(1).cloned())
1954 .or_else(|| revs.first().cloned());
1955 let (Some(from), Some(to)) = (from_rev, to_rev) else {
1956 return Err(SurfaceError::new(
1957 "target_missing",
1958 format!("no revisions to diff for {}", Json::String(doc_ref)),
1959 ));
1960 };
1961 let path: Option<String> = self
1962 .store
1963 .conn()
1964 .query_row(
1965 "SELECT path FROM docs WHERE doc_id = ?1",
1966 params![doc_id],
1967 |r| r.get(0),
1968 )
1969 .optional()?;
1970 let diff = history::diff_unified_text(&self.store, &doc_id, &from, &to)?;
1971 Ok(
1972 json!({ "doc": doc_id, "path": path.unwrap_or_default(), "from": from, "to": to, "diff": diff }),
1973 )
1974 }
1975 "docs_read_at" => {
1976 let (repo, _) = self.scope(args)?;
1977 let doc_id = self.resolve_doc_from_args(&repo, args)?;
1978 let rev = arg_string(args, "rev")?;
1979 read::docs_read_at(&self.store, &doc_id, &rev)?.ok_or_else(|| {
1980 SurfaceError::with_data(
1981 "target_missing",
1982 format!(
1983 "no revision {} for document {doc_id}",
1984 Json::String(rev.clone())
1985 ),
1986 json!({ "doc": doc_id, "rev": rev }),
1987 )
1988 })
1989 }
1990 "docs_history" => {
1991 let (repo, _) = self.scope(args)?;
1992 let glob = arg_str(args, "path_glob").filter(|g| !g.is_empty());
1993 let doc = arg_str(args, "doc").filter(|d| !d.is_empty());
1994 if glob.is_none() && doc.is_none() {
1995 return Err(SurfaceError::new(
1996 "target_missing",
1997 "docs_history requires one of path_glob or doc",
1998 ));
1999 }
2000 let include_deleted = arg_bool(args, "include_deleted")?.unwrap_or(false);
2001 if let Some(d) = doc {
2002 if history::resolve_doc_row(self.store.conn(), &repo, d, include_deleted)?
2003 .is_none()
2004 {
2005 return Err(SurfaceError::new(
2006 "doc_missing",
2007 format!("no document for {}", Json::String(d.to_owned())),
2008 ));
2009 }
2010 }
2011 history::docs_history(
2012 &self.store,
2013 &repo,
2014 glob,
2015 doc,
2016 include_deleted,
2017 arg_i64(args, "limit")?,
2018 )
2019 }
2020 "changes_since" => {
2021 let (repo, _) = self.scope(args)?;
2022 let page = self.store.changes_since(
2023 &repo,
2024 arg_i64(args, "cursor")?.unwrap_or(0),
2025 arg_usize(args, "limit")?.unwrap_or(50),
2026 arg_str(args, "origin").filter(|o| !o.is_empty()),
2027 )?;
2028 Ok(page.to_json())
2029 }
2030 "repos_status" => {
2031 let (repo, root) = self.scope(args)?;
2032 let fs = RealFileSystem;
2033 let disk = root
2034 .as_deref()
2035 .map(|r| (&fs as &dyn omgbase_sync::FileSystem, Path::new(r)));
2036 Ok(omgbase_sync::repos_status(&self.store, &repo, disk)?.to_json())
2037 }
2038 "sync_status" => {
2039 let (repo, root) = self.scope(args)?;
2040 let fs = RealFileSystem;
2041 let disk = root
2042 .as_deref()
2043 .map(|r| (&fs as &dyn omgbase_sync::FileSystem, Path::new(r)));
2044 Ok(omgbase_sync::sync_status(&self.store, &repo, disk)?.to_json())
2045 }
2046 "repos" => {
2047 let rows = self.repo_rows()?;
2048 Ok(json!({
2049 "repos": rows.iter().map(|r| json!({ "slug": r.slug, "hasSource": r.root_path.is_some() })).collect::<Vec<_>>(),
2050 }))
2051 }
2052 other => Err(SurfaceError::other(format!("unknown tool {other}"))),
2053 }
2054 }
2055}
2056
2057#[must_use]
2061pub fn apply_json(res: &ApplyResult) -> Json {
2062 let mut m = Map::new();
2063 m.insert(
2064 "results".to_owned(),
2065 Json::Array(
2066 res.results
2067 .iter()
2068 .map(|r| {
2069 let mut o = Map::new();
2070 o.insert("ids".to_owned(), json!(r.ids));
2071 if let Some(rm) = &r.removed {
2072 o.insert("removed".to_owned(), json!(rm));
2073 }
2074 if let Some(mi) = &r.merged_into {
2075 o.insert("mergedInto".to_owned(), json!(mi));
2076 }
2077 Json::Object(o)
2078 })
2079 .collect(),
2080 ),
2081 );
2082 m.insert(
2083 "revisions".to_owned(),
2084 Json::Array(
2085 res.revisions
2086 .iter()
2087 .map(|r| json!({ "doc": r.doc, "path": r.path }))
2088 .collect(),
2089 ),
2090 );
2091 if let Some(diffs) = &res.diffs {
2092 let mut d = Map::new();
2093 for (path, diff) in diffs {
2094 d.insert(
2095 path.clone(),
2096 json!({ "before": diff.before, "after": diff.after }),
2097 );
2098 }
2099 m.insert("diffs".to_owned(), Json::Object(d));
2100 }
2101 m.insert("committed".to_owned(), json!(res.committed));
2102 Json::Object(m)
2103}
2104
2105fn doc_op_json(res: &omgbase_store::DocOpResult) -> Json {
2108 let mut m = Map::new();
2109 m.insert("docId".to_owned(), json!(res.doc_id));
2110 m.insert("path".to_owned(), json!(res.path));
2111 m.insert("committed".to_owned(), json!(res.committed));
2112 if let Some(diffs) = &res.diffs {
2113 m.insert("diffs".to_owned(), diffs_json(diffs));
2114 }
2115 Json::Object(m)
2116}
2117
2118#[must_use]
2121pub fn opset_json(opset: &Opset) -> Json {
2122 let mut m = Map::new();
2123 m.insert("version".to_owned(), json!(1));
2124 m.insert("kind".to_owned(), json!("doc_update"));
2125 m.insert(
2126 "target".to_owned(),
2127 json!({ "doc": opset.target_doc, "path": opset.target_path }),
2128 );
2129 m.insert(
2130 "precondition".to_owned(),
2131 json!({
2132 "doc": opset.precondition.doc,
2133 "path": opset.precondition.path,
2134 "baseRevision": opset.precondition.base_revision,
2135 "baseContentHash": opset.precondition.base_content_hash,
2136 }),
2137 );
2138 m.insert("matcherV".to_owned(), json!(opset.matcher_v));
2139 m.insert(
2140 "ops".to_owned(),
2141 Json::Array(
2142 opset
2143 .ops
2144 .iter()
2145 .map(|p| {
2146 let mut o = Map::new();
2147 let mut op = p.op.to_json();
2152 if let Some(m) = op.as_object_mut()
2153 && let Some(c) = m.remove("child_ids")
2154 {
2155 m.insert("childIds".to_owned(), c);
2156 }
2157 o.insert("op".to_owned(), op);
2158 o.insert("disposition".to_owned(), json!(p.disposition.as_str()));
2159 o.insert("blocks".to_owned(), json!(p.blocks));
2160 o.insert("confidence".to_owned(), json!(p.confidence));
2161 o.insert("reason".to_owned(), json!(p.reason));
2162 if let Some(d) = &p.detail {
2163 o.insert("detail".to_owned(), d.clone());
2164 }
2165 Json::Object(o)
2166 })
2167 .collect(),
2168 ),
2169 );
2170 if let Some(fm) = &opset.frontmatter {
2171 m.insert("frontmatter".to_owned(), json!({ "raw": fm }));
2172 }
2173 m.insert("summary".to_owned(), opset.summary.to_json());
2174 m.insert("converges".to_owned(), json!(opset.converges));
2175 m.insert("diagnostics".to_owned(), json!(opset.diagnostics));
2176 Json::Object(m)
2177}
2178
2179fn evidence_json(e: &omgbase_store::Evidence) -> Json {
2182 let mut m = Map::new();
2183 m.insert("rrf".to_owned(), json!(e.rrf));
2184 m.insert("boosts".to_owned(), e.boosts.to_json());
2185 if let Some(r) = e.fts_rank {
2186 m.insert("ftsRank".to_owned(), json!(r));
2187 }
2188 if let Some(r) = e.vector_rank {
2189 m.insert("vectorRank".to_owned(), json!(r));
2190 if let Some(c) = e.cosine {
2191 m.insert("cosine".to_owned(), json!(c));
2192 }
2193 }
2194 Json::Object(m)
2195}
2196
2197fn observe_json(o: &omgbase_store::ObserveOutcome) -> Json {
2200 json!({
2201 "docId": o.doc_id,
2202 "path": o.path,
2203 "rev": o.rev,
2204 "commitId": o.commit_id,
2205 "converged": o.converged,
2206 "echo": o.echo,
2207 "conflicted": o.conflicted,
2208 "dispositions": o.dispositions.iter().map(|(k, n)| json!({ "kind": k, "count": n })).collect::<Vec<_>>(),
2209 })
2210}
2211
2212fn verb(d: omgbase_store::mutate_kernel::PlanDisposition) -> &'static str {
2213 use omgbase_store::mutate_kernel::PlanDisposition as D;
2214 match d {
2215 D::Same => "KEEP ",
2216 D::Edited | D::EditedMoved => "UPDATE",
2217 D::Moved => "MOVE ",
2218 D::Inserted => "INSERT",
2219 D::Deleted => "REMOVE",
2220 D::SplitFrom => "SPLIT ",
2221 D::MergedInto => "MERGE ",
2222 D::CopiedFrom => "COPY ",
2223 D::Resurrected => "RESURR",
2224 D::BulkRewrite => "REWRITE",
2225 D::Retiled => "RETILE",
2226 }
2227}
2228
2229#[must_use]
2231pub fn render_opset_plan(opset: &Opset) -> String {
2232 let mut lines = Vec::new();
2233 for p in &opset.ops {
2234 let subject = p.blocks.first().map_or("(new)", String::as_str);
2235 let conf = p.confidence.map_or(String::new(), |c| format!(" ~{c:.2}"));
2236 let why = p
2237 .reason
2238 .as_ref()
2239 .map_or(String::new(), |r| format!(" [{r}]"));
2240 lines.push(format!(
2241 "{} {:<9} {}{conf}{why}",
2242 verb(p.disposition),
2243 subject,
2244 p.disposition.as_str()
2245 ));
2246 }
2247 let s = &opset.summary;
2248 lines.push(String::new());
2249 lines.push(format!(
2250 "preserved: {} updated: {} moved: {} created: {} removed: {} split: {} merged: {} ambiguous: {}",
2251 s.preserved, s.updated, s.moved, s.created, s.removed, s.split, s.merged, s.ambiguous
2252 ));
2253 if !opset.converges {
2254 lines.push(
2255 "WARNING: plan does not reproduce the proposed content exactly — will not apply."
2256 .to_owned(),
2257 );
2258 }
2259 lines.join("\n")
2260}
2261
2262#[cfg(test)]
2263mod tests {
2264 use super::*;
2265 use omgbase_store::{MemDocStore, SequentialMinter};
2266
2267 fn surface() -> Surface {
2268 let mut store =
2269 Store::open_in_memory_with_minter(Box::new(SequentialMinter::new())).unwrap();
2270 let repo = store.create_repo("fixture").unwrap();
2271 let mut n = 0;
2272 Surface::new(store, &repo, None)
2273 .with_doc_store(Box::new(MemDocStore::new()))
2274 .with_clock(move || {
2275 n += 1;
2276 format!("2026-09-26T10:{n:02}:00.000Z")
2277 })
2278 }
2279
2280 #[test]
2281 fn catalog_lists_every_tool_of_the_table() {
2282 let names: Vec<&str> = tools().iter().map(|t| t.name).collect();
2283 for want in [
2284 "docs_outline",
2285 "docs_read",
2286 "docs_get_many",
2287 "nodes_get",
2288 "nodes_get_many",
2289 "read_ref",
2290 "docs_tree",
2291 "docs_list",
2292 "query_syntax",
2293 "query",
2294 "graph",
2295 "text_search",
2296 "resolve",
2297 "apply",
2298 "blocks_insert",
2299 "blocks_update",
2300 "blocks_move",
2301 "blocks_remove",
2302 "blocks_split",
2303 "blocks_merge",
2304 "tasks_complete",
2305 "node_set",
2306 "sections_append",
2307 "docs_append",
2308 "links_retarget",
2309 "links_stale",
2310 "links_repair",
2311 "docs_create",
2312 "docs_move",
2313 "docs_delete",
2314 "docs_set_meta",
2315 "docs_plan_update",
2316 "docs_update",
2317 "observe",
2318 "observe_many",
2319 "observe_delete",
2320 "history_node",
2321 "diff",
2322 "diff_unified",
2323 "docs_read_at",
2324 "docs_history",
2325 "changes_since",
2326 "repos_status",
2327 "sync_status",
2328 "repos",
2329 ] {
2330 assert!(names.contains(&want), "missing {want}");
2331 }
2332 assert_eq!(names.len(), 45);
2333 for t in tools() {
2334 assert_eq!(t.input_schema["type"], "object");
2335 }
2336 }
2337
2338 #[test]
2339 fn observe_read_and_query_round_trip() {
2340 let mut s = surface();
2341 let out = s.call("observe", json!({ "path": "a.md", "content": "---\nlayer: canon\n---\n# Title\n\nHello world.\n\n- [ ] task one\n" }));
2342 assert!(!out.is_error, "{}", out.body);
2343 assert_eq!(out.body["docId"], "d_0");
2344 assert_eq!(out.body["echo"], false);
2345 let read = s.call("docs_read", json!({ "doc": "a.md", "include_ids": true }));
2346 assert!(!read.is_error);
2347 assert_eq!(read.body["path"], "a.md");
2348 assert_eq!(read.body["properties"]["frontmatter"]["layer"], "canon");
2349 assert!(read.body["ids"].as_array().unwrap().len() >= 3);
2350 let q = s.call("query", json!({ "query": "select $title, layer, t: nodes collect { value where kind == \"md:task\" } from docs where layer == \"canon\" && nodes count { where kind == \"md:task\" } == 1" }));
2351 assert!(!q.is_error, "{}", q.body);
2352 assert_eq!(q.body["hits"][0]["$title"], "Title");
2353 assert_eq!(q.body["hits"][0]["layer"], "canon");
2354 assert_eq!(q.body["hits"][0]["t"][0]["value"], "task one");
2355 assert_eq!(q.body["hits"][0]["id"], "d_0");
2356 assert_eq!(q.body["consumer"], "collect");
2357 let c = s.call(
2358 "query",
2359 json!({ "query": "$repo.blocks count { where text(\"hello\") }" }),
2360 );
2361 assert_eq!(c.body["count"], 1);
2362 let bad = s.call(
2363 "query",
2364 json!({ "query": "from docs where path == \"a.md\"" }),
2365 );
2366 assert!(bad.is_error);
2367 assert_eq!(bad.body["error"], "filter_invalid");
2368 assert!(
2369 bad.body["message"]
2370 .as_str()
2371 .unwrap()
2372 .contains("did you mean the intrinsic $path")
2373 );
2374 let sem = s.call(
2375 "query",
2376 json!({ "query": "from docs where semantic(\"x\") > 0.5" }),
2377 );
2378 assert_eq!(sem.body["error"], "semantic_unavailable");
2379 let outline = s.call("docs_outline", json!({ "path": "a.md" }));
2380 let text = outline.body["text"].as_str().unwrap();
2381 assert!(text.starts_with("b_0 h1 Title §"), "{text}");
2382 assert!(text.contains("☐ task one"));
2383 let missing = s.call("docs_read", json!({ "doc": "nope.md" }));
2384 assert_eq!(missing.body["error"], "doc_missing");
2385 let repos = s.call("repos", json!({}));
2386 assert_eq!(repos.body["repos"][0]["slug"], "fixture");
2387 let unknown = s.call("docs_list", json!({ "repo": "zzz" }));
2388 assert_eq!(unknown.body["error"], "repo_not_found");
2389 }
2390
2391 #[test]
2392 fn insert_and_move_carry_the_destination_parent_cas_and_many_reads_keep_one_item() {
2393 let mut s = surface();
2394 s.call(
2395 "observe",
2396 json!({ "path": "a.md", "content": "# T\n\nOne.\n\nTwo.\n" }),
2397 );
2398 s.call("observe", json!({ "path": "b.md", "content": "# U\n" }));
2399 let stale = s.call(
2403 "blocks_insert",
2404 json!({ "to": "a.md", "markdown": "Three.", "expect": { "parent_children_hash": "00" } }),
2405 );
2406 assert!(stale.is_error, "{}", stale.body);
2407 assert_eq!(stale.body["error"], "stale_expectation");
2408 assert_eq!(stale.body["retriable"], true);
2409 let current = stale.body["data"]["current"]["parent_children_hash"]
2410 .as_str()
2411 .expect("current hash")
2412 .to_owned();
2413 assert_eq!(current.len(), 64, "{current}");
2414 let ok = s.call(
2415 "blocks_insert",
2416 json!({ "to": "a.md", "markdown": "Three.", "expect": { "parent_children_hash": current, "content_hash": "dropped, not checked" } }),
2417 );
2418 assert!(!ok.is_error, "{}", ok.body);
2419 assert_eq!(ok.body["committed"], true);
2420 let read = s.call("docs_read", json!({ "doc": "a.md" }));
2421 assert_eq!(read.body["content"], "# T\n\nOne.\n\nTwo.\n\nThree.\n");
2422 let stale = s.call(
2424 "blocks_move",
2425 json!({ "blocks": ["b_1", "b_2"], "to": "a.md", "at": "end", "expect": { "parent_children_hash": "00" } }),
2426 );
2427 assert_eq!(stale.body["error"], "stale_expectation", "{}", stale.body);
2428 let current = stale.body["data"]["current"]["parent_children_hash"]
2429 .as_str()
2430 .expect("current hash")
2431 .to_owned();
2432 let ok = s.call(
2433 "blocks_move",
2434 json!({ "blocks": ["b_1", "b_2"], "to": "a.md", "at": "end", "expect": { "parent_children_hash": current } }),
2435 );
2436 assert!(!ok.is_error, "{}", ok.body);
2437 let read = s.call("docs_read", json!({ "doc": "a.md" }));
2438 assert_eq!(read.body["content"], "# T\n\nThree.\n\nOne.\n\nTwo.\n\n");
2439 let dry = s.call(
2441 "blocks_move",
2442 json!({ "blocks": ["b_1"], "to": "a.md", "at": "start", "dry_run": true, "expect": { "parent_children_hash": "00" } }),
2443 );
2444 assert_eq!(dry.body["error"], "stale_expectation", "{}", dry.body);
2445 let many = s.call(
2448 "docs_get_many",
2449 json!({ "docs": ["a.md", "b.md"], "budget_tokens": 1 }),
2450 );
2451 assert_eq!(
2452 many.body["items"].as_array().unwrap().len(),
2453 1,
2454 "{}",
2455 many.body
2456 );
2457 assert_eq!(many.body["items"][0]["path"], "a.md");
2458 assert_eq!(many.body["truncated"], true);
2459 let one = s.call(
2460 "docs_get_many",
2461 json!({ "docs": ["b.md"], "budget_tokens": 1 }),
2462 );
2463 assert_eq!(one.body["items"].as_array().unwrap().len(), 1);
2464 assert_eq!(one.body["truncated"], false);
2465 let nodes = s.call(
2466 "nodes_get_many",
2467 json!({ "ids": ["b_zzz", "b_1", "b_2"], "resolution": "full", "budget_tokens": 1 }),
2468 );
2469 assert_eq!(
2470 nodes.body["nodes"].as_array().unwrap().len(),
2471 1,
2472 "{}",
2473 nodes.body
2474 );
2475 assert_eq!(nodes.body["nodes"][0]["id"], "b_1");
2476 assert_eq!(nodes.body["truncated"], true);
2477 assert_eq!(nodes.body["unresolved"], json!(["b_zzz"]));
2478 }
2479
2480 #[test]
2481 fn writes_go_through_the_fixed_doc_store_and_fire_the_hook() {
2482 use std::cell::Cell;
2483 use std::rc::Rc;
2484 let fired = Rc::new(Cell::new(0));
2485 let f2 = Rc::clone(&fired);
2486 let mut s = surface().with_mutation_hook(move || f2.set(f2.get() + 1));
2487 s.call(
2488 "observe",
2489 json!({ "path": "a.md", "content": "# T\n\nOne.\n" }),
2490 );
2491 assert_eq!(fired.get(), 1);
2492 let dry = s.call(
2493 "blocks_insert",
2494 json!({ "to": "a.md", "markdown": "Two.", "dry_run": true }),
2495 );
2496 assert!(!dry.is_error, "{}", dry.body);
2497 assert_eq!(dry.body["committed"], false);
2498 assert_eq!(fired.get(), 1, "a dry run never fires the hook");
2499 let wet = s.call("blocks_insert", json!({ "to": "a.md", "markdown": "Two." }));
2500 assert!(!wet.is_error, "{}", wet.body);
2501 assert_eq!(fired.get(), 2);
2502 let read = s.call("docs_read", json!({ "doc": "d_0" }));
2503 assert_eq!(read.body["content"], "# T\n\nOne.\n\nTwo.\n");
2504 let upd = s.call(
2505 "blocks_update",
2506 json!({ "block": "b_1", "markdown": "One, edited." }),
2507 );
2508 assert!(!upd.is_error, "{}", upd.body);
2509 assert_eq!(upd.body["id"], "b_1");
2510 let app = s.call(
2511 "sections_append",
2512 json!({ "heading": "T", "markdown": "Three." }),
2513 );
2514 assert!(!app.is_error, "{}", app.body);
2515 let amb = s.call(
2516 "sections_append",
2517 json!({ "heading": "Nope", "markdown": "x" }),
2518 );
2519 assert_eq!(amb.body["error"], "parent_missing");
2520 let hist = s.call("history_node", json!({ "id": "b_1" }));
2521 let entries = hist.body.as_array().unwrap();
2522 assert!(entries.len() >= 2, "{}", hist.body);
2523 assert_eq!(entries[0]["origin"], "api", "newest first");
2524 let du = s.call("diff_unified", json!({ "doc": "a.md" }));
2525 assert!(!du.is_error, "{}", du.body);
2526 assert!(du.body["diff"].as_str().unwrap().contains("+Three."));
2527 }
2528
2529 #[test]
2530 fn doc_tools_dry_run_preview_diffs_and_commit_nothing() {
2531 use std::cell::Cell;
2532 use std::rc::Rc;
2533 let fired = Rc::new(Cell::new(0));
2534 let f2 = Rc::clone(&fired);
2535 let mut s = surface().with_mutation_hook(move || f2.set(f2.get() + 1));
2536 let b_src = "# B\n\nTarget.\n";
2537 let a_src = "# A\n\nSee [b](b.md).\n";
2538 let b = s.call("docs_create", json!({ "path": "b.md", "markdown": b_src }));
2539 assert_eq!(b.body["docId"], "d_0", "{}", b.body);
2540 let a = s.call("docs_create", json!({ "path": "a.md", "markdown": a_src }));
2541 assert_eq!(a.body["docId"], "d_1", "{}", a.body);
2542 assert_eq!(fired.get(), 2);
2543 let count = |s: &Surface, table: &str| -> i64 {
2544 s.store
2545 .conn()
2546 .query_row(&format!("SELECT count(*) FROM {table}"), [], |r| r.get(0))
2547 .unwrap()
2548 };
2549 let snapshot = |s: &Surface| {
2550 (
2551 count(s, "docs"),
2552 count(s, "commits"),
2553 count(s, "revisions"),
2554 count(s, "blocks"),
2555 )
2556 };
2557 let before = snapshot(&s);
2558 let keys = |v: &Json| -> Vec<String> { v.as_object().unwrap().keys().cloned().collect() };
2559
2560 let c = s.call(
2562 "docs_create",
2563 json!({ "path": "c.md", "markdown": "# C", "frontmatter": { "title": "C" }, "dry_run": true }),
2564 );
2565 assert!(!c.is_error, "{}", c.body);
2566 assert_eq!(
2567 serde_json::to_string(&c.body).unwrap(),
2568 r#"{"docId":"d_2","path":"c.md","committed":false,"diffs":{"c.md":{"before":"","after":"---\ntitle: C\n---\n\n# C\n"}}}"#,
2569 "the reference's key order, byte for byte"
2570 );
2571
2572 let m = s.call(
2575 "docs_move",
2576 json!({ "doc": "b.md", "to_path": "notes/b.md", "retarget_inbound": true, "dry_run": true }),
2577 );
2578 assert!(!m.is_error, "{}", m.body);
2579 assert_eq!(
2583 serde_json::to_string(&m.body).unwrap(),
2584 r##"{"docId":"d_0","path":"notes/b.md","committed":false,"diffs":{"b.md":{"before":"# B\n\nTarget.\n","after":""},"notes/b.md":{"before":"","after":"# B\n\nTarget.\n"},"a.md":{"before":"# A\n\nSee [b](b.md).\n","after":"# A\n\nSee [b](notes/b.md).\n"}},"dangling":[],"retargeted":{"blocks":["b_3"],"docs":["d_1"]}}"##
2585 );
2586 assert_eq!(
2587 keys(&m.body),
2588 [
2589 "docId",
2590 "path",
2591 "committed",
2592 "diffs",
2593 "dangling",
2594 "retargeted"
2595 ]
2596 );
2597 assert_eq!(keys(&m.body["diffs"]), ["b.md", "notes/b.md", "a.md"]);
2598 assert_eq!(
2599 m.body["diffs"]["a.md"],
2600 json!({ "before": a_src, "after": "# A\n\nSee [b](notes/b.md).\n" })
2601 );
2602 let plain = s.call(
2604 "docs_move",
2605 json!({ "doc": "b.md", "to_path": "notes/b.md", "dry_run": true }),
2606 );
2607 assert_eq!(
2608 serde_json::to_string(&plain.body).unwrap(),
2609 r##"{"docId":"d_0","path":"notes/b.md","committed":false,"diffs":{"b.md":{"before":"# B\n\nTarget.\n","after":""},"notes/b.md":{"before":"","after":"# B\n\nTarget.\n"}},"dangling":[{"doc":"d_1","path":"a.md","block":"b_3","target":"b.md","anchor":null}],"retargeted":null}"##
2610 );
2611
2612 let d = s.call("docs_delete", json!({ "doc": "a.md", "dry_run": true }));
2614 assert!(!d.is_error, "{}", d.body);
2615 assert_eq!(
2616 d.body,
2617 json!({ "docId": "d_1", "path": "a.md", "committed": false, "diffs": { "a.md": { "before": a_src, "after": "" } } })
2618 );
2619 assert_eq!(keys(&d.body), ["docId", "path", "committed", "diffs"]);
2620
2621 let sm = s.call(
2623 "docs_set_meta",
2624 json!({ "doc": "a.md", "set": { "status": "open" }, "dry_run": true }),
2625 );
2626 assert!(!sm.is_error, "{}", sm.body);
2627 assert_eq!(
2628 sm.body,
2629 json!({ "docId": "d_1", "path": "a.md", "committed": false, "diffs": { "a.md": { "before": a_src, "after": "---\nstatus: open\n---\n\n# A\n\nSee [b](b.md).\n" } } })
2630 );
2631
2632 assert_eq!(fired.get(), 2, "dry runs never fire the mutation hook");
2634 assert_eq!(snapshot(&s), before);
2635 assert_eq!(
2636 s.call("docs_read", json!({ "doc": "d_0" })).body["path"],
2637 "b.md"
2638 );
2639 assert_eq!(
2640 s.call("docs_read", json!({ "doc": "d_1" })).body["content"],
2641 a_src
2642 );
2643 let taken = s.call(
2644 "docs_create",
2645 json!({ "path": "a.md", "markdown": "x", "dry_run": true }),
2646 );
2647 assert_eq!(taken.body["error"], "path_taken");
2648 let missing = s.call("docs_delete", json!({ "doc": "nope.md", "dry_run": true }));
2649 assert_eq!(missing.body["error"], "doc_missing");
2650 let real = s.call(
2652 "docs_create",
2653 json!({ "path": "c.md", "markdown": "# C\n" }),
2654 );
2655 assert_eq!(
2656 real.body,
2657 json!({ "docId": "d_3", "path": "c.md", "committed": true })
2658 );
2659 assert_eq!(fired.get(), 3);
2660 }
2661}