Skip to content

Commit 8288408

Browse files
fix(import): fail closed on truncated previews and denials
1 parent 1c8fe65 commit 8288408

4 files changed

Lines changed: 92 additions & 6 deletions

File tree

engraphis/obsidian_import.py

Lines changed: 24 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -277,7 +277,7 @@ def preview(
277277
attachment_manifest: Optional[list[dict]] = None,
278278
) -> dict:
279279
policy = self._policy(on_conflict)
280-
vault, items = self._preview_manifest(
280+
vault, items, manifest_complete = self._preview_manifest(
281281
scan, workspace_id=workspace_id, repo_id=repo_id,
282282
session_id=session_id, vault_id=vault_id, manifest=manifest,
283283
scope=scope, memory_type=memory_type, strict_root=strict_root,
@@ -286,11 +286,18 @@ def preview(
286286
plans, missing = self._plan(
287287
scan, identity, items, inspect_memories=manifest is None,
288288
)
289+
# A bounded manifest is enough to plan visible notes, but not enough to claim
290+
# that an unseen historical row was deleted. Match confirmed execution by
291+
# surfacing those rows as deferred rather than actionable ``missing`` items.
292+
pending_missing = missing if not manifest_complete else []
293+
if pending_missing:
294+
missing = []
289295
return self._report(
290296
plans, missing, scan, state="preview", vault_id=(vault or {}).get("id"),
291297
workspace_id=workspace_id, repo_id=repo_id, session_id=session_id,
292298
scope=scope, memory_type=memory_type, policy=policy,
293299
vault_label=vault_label, attachment_manifest=attachment_manifest,
300+
pending_missing=pending_missing, manifest_complete=manifest_complete,
294301
)
295302

296303
def import_scan(
@@ -426,6 +433,7 @@ def import_scan(
426433
workspace_id=workspace_id, repo_id=repo_id, session_id=session_id,
427434
scope=scope, memory_type=memory_type, policy=policy,
428435
vault_label=vault_label, attachment_manifest=attachment_manifest,
436+
manifest_complete=manifest_complete,
429437
)
430438
if terminal_state == "completed" and report["counts"].get("conflict", 0):
431439
terminal_state = "partial"
@@ -476,7 +484,7 @@ def _preview_manifest(
476484
repo_id: Optional[str], session_id: Optional[str], vault_id: Optional[str],
477485
scope: Scope, memory_type: MemoryType, manifest: Optional[dict],
478486
strict_root: bool,
479-
) -> tuple[Optional[dict], list[dict]]:
487+
) -> tuple[Optional[dict], list[dict], bool]:
480488
vaults = (
481489
list((manifest or {}).get("vaults") or [])
482490
if manifest is not None else self.store.list_source_vaults(kind=self.SOURCE_KIND)
@@ -486,6 +494,7 @@ def _preview_manifest(
486494
if manifest is not None else []
487495
)
488496
vault: Optional[dict] = None
497+
manifest_complete = True
489498
if vault_id:
490499
vault = (
491500
self.store.get_source_vault(vault_id)
@@ -513,12 +522,12 @@ def _preview_manifest(
513522
if manifest is None:
514523
# Page the manifest like import_scan does so previews on manifests
515524
# larger than one list page plan against the full item set.
516-
items, _ = self._all_source_items(vault_id=str(vault["id"]))
525+
items, manifest_complete = self._all_source_items(vault_id=str(vault["id"]))
517526
else:
518527
items = [row for row in items if row.get("vault_id") == vault.get("id")]
519528
else:
520529
items = []
521-
return vault, items
530+
return vault, items, manifest_complete
522531

523532
def _resolve_or_register_vault(
524533
self, scan: _ImportScan, *, workspace_id: str,
@@ -1332,6 +1341,8 @@ def _report(
13321341
repo_id: Optional[str], session_id: Optional[str], scope: Scope,
13331342
memory_type: MemoryType, policy: str, vault_label: str,
13341343
attachment_manifest: Optional[list[dict]],
1344+
pending_missing: Optional[list[dict]] = None,
1345+
manifest_complete: bool = True,
13351346
) -> dict:
13361347
files = [self._preview_row(plan) for plan in plans]
13371348
files.extend({
@@ -1346,11 +1357,16 @@ def _report(
13461357
"relative_path": str(item.get("relative_path") or ""), "status": "missing",
13471358
"action": "missing", "reason": "source_not_seen", "warnings": [],
13481359
} for item in missing)
1360+
files.extend({
1361+
"relative_path": str(item.get("relative_path") or ""), "status": "pending",
1362+
"action": "pending", "reason": "missing_check_deferred", "warnings": [],
1363+
} for item in pending_missing or [])
13491364
return self._report_payload(
13501365
files, scan, state=state, vault_id=vault_id,
13511366
workspace_id=workspace_id, repo_id=repo_id, session_id=session_id,
13521367
scope=scope, memory_type=memory_type, policy=policy,
13531368
vault_label=vault_label, attachment_manifest=attachment_manifest,
1369+
manifest_complete=manifest_complete,
13541370
)
13551371

13561372
def _final_report(
@@ -1360,6 +1376,7 @@ def _final_report(
13601376
import_id: str, workspace_id: str, repo_id: Optional[str],
13611377
session_id: Optional[str], scope: Scope, memory_type: MemoryType,
13621378
policy: str, vault_label: str, attachment_manifest: Optional[list[dict]],
1379+
manifest_complete: bool = True,
13631380
) -> dict:
13641381
files = list(outcomes)
13651382
processed_paths = {
@@ -1396,6 +1413,7 @@ def _final_report(
13961413
workspace_id=workspace_id, repo_id=repo_id, session_id=session_id,
13971414
scope=scope, memory_type=memory_type, policy=policy,
13981415
vault_label=vault_label, attachment_manifest=attachment_manifest,
1416+
manifest_complete=manifest_complete,
13991417
)
14001418
report.update({"job_id": job_id, "import_id": import_id})
14011419
return report
@@ -1423,6 +1441,7 @@ def _report_payload(
14231441
vault_id: Optional[str], workspace_id: Optional[str], repo_id: Optional[str],
14241442
session_id: Optional[str], scope: Scope, memory_type: MemoryType,
14251443
policy: str, vault_label: str, attachment_manifest: Optional[list[dict]],
1444+
manifest_complete: bool = True,
14261445
) -> dict:
14271446
counts: dict[str, int] = {}
14281447
for row in files:
@@ -1440,6 +1459,7 @@ def _report_payload(
14401459
formats[name] = formats.get(name, 0) + 1
14411460
return {
14421461
"state": state, "status": state, "vault_id": vault_id,
1462+
"manifest_complete": bool(manifest_complete),
14431463
"vault_label": str(vault_label or "")[:200],
14441464
"source_id": vault_id,
14451465
"source_label": str(vault_label or "")[:200],

engraphis/routes/v2_api.py

Lines changed: 7 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3126,10 +3126,16 @@ def _mark_authoritative_denial() -> None:
31263126
global _authoritative_denial_at, _denied_state_digests
31273127
with _ENTITLEMENT_REFRESH_LOCK:
31283128
_authoritative_denial_at = time.time()
3129+
# Publish the fail-closed state before probing either persisted source. Those
3130+
# reads can block on a slow state directory; a concurrent license/bootstrap
3131+
# request must not keep serving the pre-denial grants while they are in flight.
3132+
# Keep the baseline empty until both probes complete so an active-looking record
3133+
# cannot be mistaken for a post-denial reconnect during the capture window.
3134+
_AUTHORITATIVE_DENIAL_PENDING.set()
3135+
_denied_state_digests = {}
31293136
_denied_state_digests = {
31303137
source: _persisted_state_digest(source) for source in ("session", "cloud")
31313138
}
3132-
_AUTHORITATIVE_DENIAL_PENDING.set()
31333139

31343140

31353141
def _clear_superseded_denial(known_source: str, observed_digest: Optional[str]) -> bool:

tests/test_hosted_plan_resolution.py

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1134,6 +1134,32 @@ def _blocked_write():
11341134
assert not worker.is_alive()
11351135

11361136

1137+
def test_denial_guard_publishes_before_blocked_digest_probe(monkeypatch) -> None:
1138+
"""Readers fail closed while the pre-persistence denial baseline is captured."""
1139+
1140+
_connect(monkeypatch, pinned_token=False)
1141+
entered = threading.Event()
1142+
release = threading.Event()
1143+
1144+
def _blocked_digest(_source):
1145+
entered.set()
1146+
assert v2_api._AUTHORITATIVE_DENIAL_PENDING.is_set()
1147+
assert v2_api._denied_state_digests == {}
1148+
assert release.wait(timeout=5.0)
1149+
return "before-denial"
1150+
1151+
monkeypatch.setattr(v2_api, "_persisted_state_digest", _blocked_digest)
1152+
worker = threading.Thread(target=v2_api._mark_authoritative_denial)
1153+
worker.start()
1154+
assert entered.wait(timeout=5.0)
1155+
try:
1156+
assert v2_api._AUTHORITATIVE_DENIAL_PENDING.is_set()
1157+
finally:
1158+
release.set()
1159+
worker.join(timeout=5.0)
1160+
assert not worker.is_alive()
1161+
1162+
11371163
def test_a_transport_failure_is_not_mistaken_for_a_billing_denial(monkeypatch) -> None:
11381164
"""Only an authoritative 401/402/403 clears access; an outage must not."""
11391165

tests/test_obsidian_service.py

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77
import pytest
88

99
from engraphis.service import MemoryService, ValidationError
10-
from engraphis.obsidian_import import scan_obsidian_upload
10+
from engraphis.obsidian_import import ObsidianImporter, scan_obsidian_upload
1111

1212

1313
_TERMINAL_STATES = {"completed", "partial", "failed", "cancelled"}
@@ -90,6 +90,40 @@ def test_preview_pages_the_full_manifest_like_execution():
9090
service.close()
9191

9292

93+
def test_preview_defers_missing_rows_when_manifest_is_truncated(monkeypatch):
94+
service = _service()
95+
try:
96+
started, imported = _import(service, [("A.md", b"# A\n")])
97+
assert imported["state"] == "completed"
98+
vault_id = started["vault_id"]
99+
service.store.upsert_source_import_item(
100+
vault_id=vault_id,
101+
source_key=hashlib.sha256(b"gone").hexdigest(),
102+
relative_path="gone.md",
103+
)
104+
items = service.store.list_source_import_items(vault_id=vault_id)
105+
106+
def _truncated(_self, *, vault_id, states=None):
107+
del vault_id, states
108+
return items, False
109+
110+
monkeypatch.setattr(ObsidianImporter, "_all_source_items", _truncated)
111+
preview = service.preview_obsidian_upload(
112+
files=[("A.md", b"# A\n")], attachment_manifest=[],
113+
workspace="alpha", vault_label="Team notes", vault_id=vault_id,
114+
)
115+
116+
statuses = {
117+
row["relative_path"]: row["status"] for row in preview["files"]
118+
}
119+
assert preview["manifest_complete"] is False
120+
assert preview["counts"].get("missing", 0) == 0
121+
assert preview["counts"]["pending"] == 1
122+
assert statuses["gone.md"] == "pending"
123+
finally:
124+
service.close()
125+
126+
93127
def test_preview_is_write_free_and_service_enforces_confirmation_and_upload_guards(monkeypatch):
94128
service = _service()
95129
try:

0 commit comments

Comments
 (0)