diff --git a/solstone/apps/import/routes.py b/solstone/apps/import/routes.py index d16a2dfcf..1621c488d 100644 --- a/solstone/apps/import/routes.py +++ b/solstone/apps/import/routes.py @@ -25,7 +25,7 @@ from solstone.convey.reasons import ( JOURNAL_SOURCE_PROBLEM, MISSING_REQUIRED_FIELD, ) -from solstone.convey.utils import error_response, respond_collection +from solstone.convey.utils import error_response, load_json, respond_collection from solstone.think.detect_created import detect_created from solstone.think.importers.utils import ( build_import_info, @@ -1001,6 +1001,74 @@ def api_journal_source_status(name: str) -> Any: ) +@import_bp.route("/api/journal-sources//staged") +def api_journal_source_staged(name: str) -> Any: + source = find_journal_source_by_name(name) + if not source: + return error_response( + JOURNAL_SOURCE_PROBLEM, + status=404, + detail=f"Journal source '{name}' not found", + ) + area = request.args.get("area") + if area is not None and area not in {"entities", "facets", "config"}: + return error_response( + INVALID_REQUEST_VALUE, + status=400, + detail="Area must be one of: entities, facets, config", + ) + # Mirrors api_journal_source_status: a registry-returned record always has a + # valid prefix, so call journal_source_state_prefix directly (no try/except). + state_dir = get_state_directory(journal_source_state_prefix(source)) + items: list[dict[str, Any]] = [] + + if area in {None, "entities"}: + staged_dir = state_dir / "entities" / "staged" + for staged_path in sorted(staged_dir.glob("*.json")): + payload = load_json(staged_path) + if not isinstance(payload, dict): + continue + items.append( + { + "area": "entities", + "source_id": staged_path.stem, + "reason": payload.get("reason"), + "source_entity": payload.get("source_entity"), + "match_candidates": payload.get("match_candidates"), + "staged_at": payload.get("staged_at"), + } + ) + + if area in {None, "facets"}: + staged_dir = state_dir / "facets" / "staged" + for staged_path in sorted(staged_dir.glob("**/*.staged.json")): + payload = load_json(staged_path) + if not isinstance(payload, dict): + continue + relative_path = staged_path.relative_to(staged_dir) + parts = relative_path.parts + if len(parts) < 3: + continue + line = { + "area": "facets", + "staged_file": relative_path.as_posix(), + "facet": parts[0], + "file_type": parts[1], + } + line.update(payload) + items.append(line) + + if area in {None, "config"}: + diff = load_json(state_dir / "config" / "diff.json") + # Config-parity decision: include the config item iff diff.json exists + # AND loads as a dict. A missing / unreadable / non-dict diff is omitted + # (still HTTP 200) — never a 500, never an empty-diff placeholder. + if isinstance(diff, dict): + items.append({"area": "config", "diff": diff}) + + return respond_collection(items) + + @import_bp.route("/journal//manifest/") @require_journal_source def journal_source_manifest(key_prefix: str, area: str) -> Any: diff --git a/tests/test_journal_source_dl_surfaces.py b/tests/test_journal_source_dl_surfaces.py index da2928c77..7f66d306b 100644 --- a/tests/test_journal_source_dl_surfaces.py +++ b/tests/test_journal_source_dl_surfaces.py @@ -4,6 +4,7 @@ from __future__ import annotations import argparse +import json from importlib import import_module import pytest @@ -16,6 +17,7 @@ journal_sources = import_module("solstone.apps.import.journal_sources") import_routes = import_module("solstone.apps.import.routes") journal_source_cli = import_module("solstone.think.importers.journal_source_cli") +create_state_directory = journal_sources.create_state_directory generate_key = journal_sources.generate_key journal_source_state_prefix = journal_sources.journal_source_state_prefix load_journal_source_by_fingerprint = journal_sources.load_journal_source_by_fingerprint @@ -80,6 +82,72 @@ def _save_dl_and_pl() -> dict: return dl_source +def _client(): + app = Flask(__name__) + app.register_blueprint(import_routes.import_bp) + return app.test_client() + + +def _write_json(path, payload) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(json.dumps(payload), encoding="utf-8") + + +def _stage_all_areas(journal_root) -> dict: + dl_source = _dl_source() + assert save_journal_source(dl_source) is True + prefix = journal_source_state_prefix(dl_source) + state_dir = create_state_directory(journal_root, prefix) + + entity_item = { + "area": "entities", + "source_id": "ent-1", + "reason": "new", + "source_entity": {"name": "Ada"}, + "match_candidates": [{"entity_id": "ada"}], + "staged_at": "2026-06-07T00:00:00Z", + } + _write_json( + state_dir / "entities" / "staged" / "ent-1.json", + { + "reason": entity_item["reason"], + "source_entity": entity_item["source_entity"], + "match_candidates": entity_item["match_candidates"], + "staged_at": entity_item["staged_at"], + "junk": "ignored", + }, + ) + _write_json(state_dir / "entities" / "staged" / "bad.json", []) + + facet_payload = {"source_id": "foo", "reason": "missing_target"} + facet_item = { + "area": "facets", + "staged_file": "work/entity_observations/foo.staged.json", + "facet": "work", + "file_type": "entity_observations", + **facet_payload, + } + _write_json( + state_dir + / "facets" + / "staged" + / "work" + / "entity_observations" + / "foo.staged.json", + facet_payload, + ) + _write_json(state_dir / "facets" / "staged" / "shallow.staged.json", {}) + + config_item = {"area": "config", "diff": {"field.a": {"category": "x"}}} + _write_json(state_dir / "config" / "diff.json", config_item["diff"]) + + return { + "entities": entity_item, + "facets": facet_item, + "config": config_item, + } + + def test_api_journal_source_list_excludes_pl_records(journal_env) -> None: dl_source = _save_dl_and_pl() app = Flask(__name__) @@ -123,3 +191,81 @@ def test_cli_status_and_revoke_cannot_target_pl_fingerprint_by_name( pl_record = load_journal_source_by_fingerprint(FINGERPRINT) assert pl_record is not None assert pl_record["revoked"] is False + + +def test_api_journal_source_staged_returns_all_areas_and_skips_invalid( + journal_env, +) -> None: + expected = _stage_all_areas(journal_env) + + response = _client().get("/app/import/api/journal-sources/alpha/staged") + + assert response.status_code == 200 + body = response.get_json() + assert body["total"] == len(body["items"]) + assert body["items"] == [ + expected["entities"], + expected["facets"], + expected["config"], + ] + assert all(item.get("source_id") != "bad" for item in body["items"]) + assert all( + item.get("staged_file") != "shallow.staged.json" for item in body["items"] + ) + + +def test_api_journal_source_staged_area_filter(journal_env) -> None: + expected = _stage_all_areas(journal_env) + client = _client() + + entities = client.get("/app/import/api/journal-sources/alpha/staged?area=entities") + facets = client.get("/app/import/api/journal-sources/alpha/staged?area=facets") + config = client.get("/app/import/api/journal-sources/alpha/staged?area=config") + + assert entities.status_code == 200 + assert entities.get_json() == {"items": [expected["entities"]], "total": 1} + assert facets.status_code == 200 + assert facets.get_json() == {"items": [expected["facets"]], "total": 1} + assert config.status_code == 200 + assert config.get_json() == {"items": [expected["config"]], "total": 1} + + +def test_api_journal_source_staged_empty_valid_and_unknown(journal_env) -> None: + dl_source = _dl_source() + assert save_journal_source(dl_source) is True + create_state_directory(journal_env, journal_source_state_prefix(dl_source)) + client = _client() + + empty = client.get("/app/import/api/journal-sources/alpha/staged") + unknown = client.get("/app/import/api/journal-sources/ghost/staged") + + assert empty.status_code == 200 + assert empty.get_json() == {"items": [], "total": 0} + assert unknown.status_code == 404 + assert unknown.get_json()["reason_code"] == "journal_source_problem" + + +def test_api_journal_source_staged_invalid_area_returns_error(journal_env) -> None: + dl_source = _dl_source() + assert save_journal_source(dl_source) is True + + response = _client().get("/app/import/api/journal-sources/alpha/staged?area=bogus") + + assert response.status_code == 400 + body = response.get_json() + assert body["reason_code"] == "invalid_request_value" + assert "items" not in body + + +def test_api_journal_source_staged_omits_non_dict_config(journal_env) -> None: + dl_source = _dl_source() + assert save_journal_source(dl_source) is True + state_dir = create_state_directory( + journal_env, journal_source_state_prefix(dl_source) + ) + _write_json(state_dir / "config" / "diff.json", []) + + response = _client().get("/app/import/api/journal-sources/alpha/staged?area=config") + + assert response.status_code == 200 + assert response.get_json() == {"items": [], "total": 0}