diff --git a/apps/import/ingest.py b/apps/import/ingest.py index ada4508d1..84fecf29a 100644 --- a/apps/import/ingest.py +++ b/apps/import/ingest.py @@ -5,14 +5,15 @@ from __future__ import annotations -import json import hashlib +import json import logging import os import re import tempfile from datetime import datetime, timezone from pathlib import Path +from typing import Any from flask import abort, g, jsonify, request from werkzeug.utils import secure_filename @@ -43,6 +44,30 @@ _DAY_RE = re.compile(r"^\d{8}$") _SEGMENT_RE = re.compile(r"^\d{6}_\d+$") _STREAM_RE = re.compile(r"^[a-z0-9][a-z0-9._-]*$") _FACET_NAME_RE = re.compile(r"^[a-z0-9][a-z0-9_-]*$") +_IMPORT_ID_RE = re.compile(r"^\d{8}_\d{6}$") + +_NEVER_TRANSFER_PATHS = frozenset( + { + "convey.password_hash", + "convey.secret", + "setup.completed_at", + "providers.auth", + "providers.key_validation", + "transcribe.whisper.device", + } +) +_NEVER_TRANSFER_PREFIXES = ("env.",) +_IDENTITY_PATHS = frozenset( + { + "identity.name", + "identity.preferred", + "identity.bio", + "identity.pronouns", + "identity.aliases", + "identity.email_addresses", + "identity.timezone", + } +) def _append_decision(log_path: Path, entry: dict) -> None: @@ -63,6 +88,31 @@ def _write_state_atomic(state_path: Path, state_data: dict) -> None: raise +def _flatten_config(cfg: dict, prefix: str = "") -> dict[str, Any]: + """Flatten a nested config dict to dot-separated paths.""" + result: dict[str, Any] = {} + for key, value in cfg.items(): + path = f"{prefix}{key}" if prefix else key + if isinstance(value, dict): + result.update(_flatten_config(value, f"{path}.")) + else: + result[path] = value + return result + + +def _is_never_transfer(path: str) -> bool: + if path in _NEVER_TRANSFER_PATHS: + return True + return any(path.startswith(prefix) for prefix in _NEVER_TRANSFER_PREFIXES) + + +def _categorize_field(path: str) -> str: + """Return category for a config field path.""" + if path in _IDENTITY_PATHS: + return "transferable" + return "preference" + + from .facet_ingest import process_facet @@ -701,3 +751,267 @@ def register_ingest_routes(bp) -> None: "errors": errors, } ) + + @bp.route("/journal//ingest/imports", methods=["POST"]) + @require_journal_source + def ingest_imports(key_prefix: str): + if g.journal_source["key"][:8] != key_prefix: + abort(403, description="Key prefix mismatch") + + payload = request.get_json(silent=True) + if not isinstance(payload, dict): + return jsonify({"error": "Invalid JSON body"}), 400 + + imports = payload.get("imports") + if not isinstance(imports, list): + return jsonify({"error": "Missing imports array"}), 400 + + state_dir = get_state_directory(key_prefix) + log_path = state_dir / "imports" / "log.jsonl" + state_path = state_dir / "imports" / "state.json" + staged_dir = state_dir / "imports" / "staged" + + try: + imports_state = json.loads(state_path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + imports_state = {} + if not isinstance(imports_state, dict): + imports_state = {} + received = imports_state.get("received") + if not isinstance(received, dict): + received = {} + imports_state = {"received": dict(received)} + + journal_root = Path(state.journal_root) + + copied = 0 + skipped = 0 + staged = 0 + errors: list[dict[str, str]] = [] + + for item in imports: + try: + if not isinstance(item, dict): + raise ValueError("Import item must be an object") + + import_id = str(item.get("id", "")).strip() + if not import_id or not _IMPORT_ID_RE.match(import_id): + raise ValueError(f"Invalid import id: {import_id!r}") + + import_json = item.get("import_json") + imported_json = item.get("imported_json") + content_manifest = item.get("content_manifest") + + if not isinstance(import_json, dict): + raise ValueError("import_json must be an object") + if not isinstance(imported_json, dict): + raise ValueError("imported_json must be an object") + if not isinstance(content_manifest, list): + raise ValueError("content_manifest must be an array") + + hash_input = json.dumps( + { + "import_json": import_json, + "imported_json": imported_json, + "content_manifest": content_manifest, + }, + sort_keys=True, + ensure_ascii=False, + ).encode() + content_hash = hashlib.sha256(hash_input).hexdigest() + + if imports_state["received"].get(import_id) == content_hash: + skipped += 1 + _append_decision( + log_path, + { + "ts": datetime.now(timezone.utc).isoformat(), + "action": "skipped", + "item_type": "import", + "item_id": import_id, + "reason": "idempotent", + }, + ) + continue + + target_dir = journal_root / "imports" / import_id + if target_dir.is_dir(): + staged_dir.mkdir(parents=True, exist_ok=True) + staged_payload = { + "import_id": import_id, + "import_json": import_json, + "imported_json": imported_json, + "content_manifest": content_manifest, + "reason": "id_collision", + "staged_at": datetime.now(timezone.utc).isoformat(), + } + (staged_dir / f"{import_id}.json").write_text( + json.dumps(staged_payload, indent=2, ensure_ascii=False) + "\n", + encoding="utf-8", + ) + staged += 1 + _append_decision( + log_path, + { + "ts": datetime.now(timezone.utc).isoformat(), + "action": "staged", + "item_type": "import", + "item_id": import_id, + "reason": "id_collision", + }, + ) + else: + target_dir.mkdir(parents=True, exist_ok=True) + (target_dir / "import.json").write_text( + json.dumps(import_json, indent=2, ensure_ascii=False) + "\n", + encoding="utf-8", + ) + (target_dir / "imported.json").write_text( + json.dumps(imported_json, indent=2, ensure_ascii=False) + "\n", + encoding="utf-8", + ) + lines = [ + json.dumps(entry, ensure_ascii=False) + for entry in content_manifest + ] + (target_dir / "content_manifest.jsonl").write_text( + "\n".join(lines) + "\n" if lines else "", + encoding="utf-8", + ) + copied += 1 + _append_decision( + log_path, + { + "ts": datetime.now(timezone.utc).isoformat(), + "action": "copied", + "item_type": "import", + "item_id": import_id, + "reason": "new", + }, + ) + + imports_state["received"][import_id] = content_hash + except Exception as exc: + import_id_str = item.get("id", "") if isinstance(item, dict) else "" + errors.append({"import_id": str(import_id_str), "error": str(exc)}) + + _write_state_atomic(state_path, imports_state) + + if copied > 0: + source = g.journal_source + source.setdefault("stats", {}) + source["stats"]["imports_received"] = ( + source["stats"].get("imports_received", 0) + copied + ) + save_journal_source(source) + + return jsonify( + { + "copied": copied, + "skipped": skipped, + "staged": staged, + "errors": errors, + } + ) + + @bp.route("/journal//ingest/config", methods=["POST"]) + @require_journal_source + def ingest_config(key_prefix: str): + if g.journal_source["key"][:8] != key_prefix: + abort(403, description="Key prefix mismatch") + + payload = request.get_json(silent=True) + if not isinstance(payload, dict): + return jsonify({"error": "Invalid JSON body"}), 400 + + source_config = payload.get("config") + if not isinstance(source_config, dict): + return jsonify({"error": "Missing config object"}), 400 + + state_dir = get_state_directory(key_prefix) + log_path = state_dir / "config" / "log.jsonl" + state_path = state_dir / "config" / "state.json" + config_dir = state_dir / "config" + + try: + config_state = json.loads(state_path.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + config_state = {} + if not isinstance(config_state, dict): + config_state = {} + + content_hash = hashlib.sha256( + json.dumps(source_config, sort_keys=True, ensure_ascii=False).encode() + ).hexdigest() + + if config_state.get("last_hash") == content_hash: + _append_decision( + log_path, + { + "ts": datetime.now(timezone.utc).isoformat(), + "action": "skipped", + "item_type": "config", + "item_id": "journal.json", + "reason": "idempotent", + }, + ) + return jsonify({"staged": False, "skipped": True, "reason": "idempotent"}) + + from think.utils import get_config + + target_config = get_config() + source_flat = _flatten_config(source_config) + target_flat = _flatten_config(target_config) + + all_keys = sorted(set(source_flat) | set(target_flat)) + diff = {} + for key in all_keys: + if _is_never_transfer(key): + continue + source_val = source_flat.get(key) + target_val = target_flat.get(key) + if source_val != target_val: + diff[key] = { + "source": source_val, + "target": target_val, + "category": _categorize_field(key), + } + + config_dir.mkdir(parents=True, exist_ok=True) + (config_dir / "source_config.json").write_text( + json.dumps(source_config, indent=2, ensure_ascii=False) + "\n", + encoding="utf-8", + ) + (config_dir / "diff.json").write_text( + json.dumps(diff, indent=2, ensure_ascii=False) + "\n", + encoding="utf-8", + ) + + config_state["last_hash"] = content_hash + _write_state_atomic(state_path, config_state) + + _append_decision( + log_path, + { + "ts": datetime.now(timezone.utc).isoformat(), + "action": "staged", + "item_type": "config", + "item_id": "journal.json", + "reason": "config_received", + }, + ) + + source = g.journal_source + source.setdefault("stats", {}) + source["stats"]["config_received"] = ( + source["stats"].get("config_received", 0) + 1 + ) + save_journal_source(source) + + return jsonify( + { + "staged": True, + "skipped": False, + "diff_fields": len(diff), + } + ) diff --git a/observe/export.py b/observe/export.py index f4fbb30e5..e773ce34c 100644 --- a/observe/export.py +++ b/observe/export.py @@ -27,8 +27,9 @@ from observe.transfer import ( _normalize_url, _parse_day_spec, ) -from think.utils import get_journal, iter_segments, setup_cli from think.entities.journal import load_all_journal_entities +from think.importers.sync import SYNCABLE_REGISTRY +from think.utils import get_config, get_journal, iter_segments, setup_cli logger = logging.getLogger(__name__) @@ -37,6 +38,17 @@ _FACET_NAME_RE = re.compile(r"^[a-z0-9][a-z0-9_-]*$") _DAY_RE = re.compile(r"^\d{8}$") _DAY_JSONL_RE = re.compile(r"^\d{8}\.jsonl$") _DAY_MD_RE = re.compile(r"^\d{8}\.md$") +_IMPORT_ID_RE = re.compile(r"^\d{8}_\d{6}$") +_NEVER_TRANSFER_PATHS = frozenset( + { + "convey.password_hash", + "convey.secret", + "setup.completed_at", + "providers.auth", + "providers.key_validation", + "transcribe.whisper.device", + } +) def _query_manifest( @@ -190,6 +202,26 @@ def _classify_facet_file(relative: PurePosixPath) -> str | None: return None +def _strip_never_transfer(config: dict) -> dict: + """Return a deep copy of config with never-transfer fields removed.""" + import copy as _copy + + result = _copy.deepcopy(config) + for path in _NEVER_TRANSFER_PATHS: + parts = path.split(".") + obj = result + for part in parts[:-1]: + if isinstance(obj, dict) and part in obj: + obj = obj[part] + else: + break + else: + if isinstance(obj, dict): + obj.pop(parts[-1], None) + result.pop("env", None) + return result + + def export_segments(base_url: str, key: str, days: list[str], dry_run: bool) -> None: session = requests.Session() session.headers["Authorization"] = f"Bearer {key}" @@ -613,6 +645,232 @@ def export_facets(base_url: str, key: str, dry_run: bool) -> None: session.close() +def export_imports(base_url: str, key: str, dry_run: bool) -> None: + """Export import metadata to a remote solstone instance.""" + session = requests.Session() + session.headers["Authorization"] = f"Bearer {key}" + + try: + try: + remote_manifest = _query_manifest(session, base_url, key, area="imports") + except requests.ConnectionError: + print(f"Connection failed: could not reach {base_url}") + return + except ValueError as e: + print(str(e)) + return + + received = remote_manifest.get("received", {}) + + journal_root = Path(get_journal()) + imports_dir = journal_root / "imports" + if not imports_dir.is_dir(): + print("No imports directory found") + return + + sync_state_names = {f"{name}.json" for name in SYNCABLE_REGISTRY} + + to_send = [] + new_count = 0 + changed_count = 0 + unchanged_count = 0 + + for entry in sorted(imports_dir.iterdir()): + if entry.is_file() and entry.name in sync_state_names: + continue + if not entry.is_dir(): + continue + if not _IMPORT_ID_RE.match(entry.name): + continue + + import_id = entry.name + import_json_path = entry / "import.json" + imported_json_path = entry / "imported.json" + manifest_path = entry / "content_manifest.jsonl" + + if not import_json_path.exists() or not imported_json_path.exists(): + continue + + try: + import_json = json.loads(import_json_path.read_text(encoding="utf-8")) + imported_json = json.loads( + imported_json_path.read_text(encoding="utf-8") + ) + content_manifest = [] + if manifest_path.exists(): + for line in manifest_path.read_text(encoding="utf-8").splitlines(): + line = line.strip() + if line: + content_manifest.append(json.loads(line)) + except (json.JSONDecodeError, OSError) as exc: + logger.warning("Failed to read import %s: %s", import_id, exc) + continue + + hash_input = json.dumps( + { + "import_json": import_json, + "imported_json": imported_json, + "content_manifest": content_manifest, + }, + sort_keys=True, + ensure_ascii=False, + ).encode() + content_hash = hashlib.sha256(hash_input).hexdigest() + + if received.get(import_id) == content_hash: + unchanged_count += 1 + continue + + if import_id in received: + changed_count += 1 + else: + new_count += 1 + + to_send.append( + { + "id": import_id, + "import_json": import_json, + "imported_json": imported_json, + "content_manifest": content_manifest, + } + ) + + if dry_run: + print( + f"Dry run: {new_count} new, {changed_count} changed, " + f"{unchanged_count} unchanged" + ) + return + + if not to_send: + print("Nothing to send - remote imports are up to date") + return + + key_prefix = key[:8] + url = f"{base_url}/app/import/journal/{key_prefix}/ingest/imports" + for attempt, delay in enumerate(RETRY_BACKOFF): + try: + response = session.post( + url, json={"imports": to_send}, timeout=UPLOAD_TIMEOUT + ) + if response.status_code == 200: + break + if response.status_code == 401: + print("Authentication failed: invalid or missing API key") + return + if response.status_code == 403: + print("Authentication failed: journal source revoked or disabled") + return + if 500 <= response.status_code <= 599: + logger.warning( + "Import upload attempt %s failed: %s %s", + attempt + 1, + response.status_code, + response.text, + ) + else: + print( + f"Import upload failed: {response.status_code} {response.text}" + ) + return + except (requests.RequestException, OSError) as e: + logger.warning("Import upload attempt %s failed: %s", attempt + 1, e) + if attempt < len(RETRY_BACKOFF) - 1: + time.sleep(delay) + else: + print("Import upload failed after all retries") + return + + result = response.json() + errors = result.get("errors", []) + if errors: + for err in errors: + print(f" Error: {err}") + print( + f"\nExport complete: {result.get('copied', 0)} copied, " + f"{result.get('staged', 0)} staged, " + f"{result.get('skipped', 0)} skipped" + ) + if errors: + print(f" {len(errors)} error(s)") + finally: + session.close() + + +def export_config(base_url: str, key: str, dry_run: bool) -> None: + """Export config snapshot to a remote solstone instance.""" + session = requests.Session() + session.headers["Authorization"] = f"Bearer {key}" + + try: + try: + remote_manifest = _query_manifest(session, base_url, key, area="config") + except requests.ConnectionError: + print(f"Connection failed: could not reach {base_url}") + return + except ValueError as e: + print(str(e)) + return + + config = _strip_never_transfer(get_config()) + content_hash = hashlib.sha256( + json.dumps(config, sort_keys=True, ensure_ascii=False).encode() + ).hexdigest() + + if remote_manifest.get("last_hash") == content_hash: + print("Nothing to send - remote config is up to date") + return + + if dry_run: + print("Dry run: config has changed, would send snapshot") + return + + key_prefix = key[:8] + url = f"{base_url}/app/import/journal/{key_prefix}/ingest/config" + for attempt, delay in enumerate(RETRY_BACKOFF): + try: + response = session.post( + url, json={"config": config}, timeout=UPLOAD_TIMEOUT + ) + if response.status_code == 200: + break + if response.status_code == 401: + print("Authentication failed: invalid or missing API key") + return + if response.status_code == 403: + print("Authentication failed: journal source revoked or disabled") + return + if 500 <= response.status_code <= 599: + logger.warning( + "Config upload attempt %s failed: %s %s", + attempt + 1, + response.status_code, + response.text, + ) + else: + print( + f"Config upload failed: {response.status_code} {response.text}" + ) + return + except (requests.RequestException, OSError) as e: + logger.warning("Config upload attempt %s failed: %s", attempt + 1, e) + if attempt < len(RETRY_BACKOFF) - 1: + time.sleep(delay) + else: + print("Config upload failed after all retries") + return + + result = response.json() + if result.get("staged"): + print( + f"\nExport complete: config staged ({result.get('diff_fields', 0)} fields differ)" + ) + elif result.get("skipped"): + print("Nothing to send - remote config is up to date") + finally: + session.close() + + def main() -> None: parser = argparse.ArgumentParser( description="Export journal data to a remote solstone instance" @@ -630,7 +888,7 @@ def main() -> None: parser.add_argument( "--only", default=None, - help="Export only specific area (segments, entities, facets)", + help="Export only specific area (segments, entities, facets, imports, config)", ) parser.add_argument( "--dry-run", @@ -654,6 +912,14 @@ def main() -> None: export_facets(base_url, args.key, args.dry_run) return + if args.only == "imports": + export_imports(base_url, args.key, args.dry_run) + return + + if args.only == "config": + export_config(base_url, args.key, args.dry_run) + return + if args.only is not None and args.only != "segments": print(f"Export of '{args.only}' is not yet implemented") sys.exit(0) @@ -663,3 +929,5 @@ def main() -> None: if args.only is None: export_entities(base_url, args.key, args.dry_run) export_facets(base_url, args.key, args.dry_run) + export_imports(base_url, args.key, args.dry_run) + export_config(base_url, args.key, args.dry_run) diff --git a/tests/test_config_ingest.py b/tests/test_config_ingest.py new file mode 100644 index 000000000..4e294b6c9 --- /dev/null +++ b/tests/test_config_ingest.py @@ -0,0 +1,243 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +from importlib import import_module + +import pytest +from flask import Blueprint, Flask + +import convey.state +import think.utils + +journal_sources = import_module("apps.import.journal_sources") +ingest = import_module("apps.import.ingest") + +create_state_directory = journal_sources.create_state_directory +generate_key = journal_sources.generate_key +get_state_directory = journal_sources.get_state_directory +load_journal_source = journal_sources.load_journal_source +save_journal_source = journal_sources.save_journal_source +register_ingest_routes = ingest.register_ingest_routes + + +@pytest.fixture +def journal_env(tmp_path, monkeypatch): + monkeypatch.setattr(convey.state, "journal_root", str(tmp_path), raising=False) + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + think.utils._journal_path_cache = None + (tmp_path / "apps" / "import" / "journal_sources").mkdir( + parents=True, exist_ok=True + ) + return tmp_path + + +def _source(name="test-source", key=None, **overrides): + if key is None: + key = generate_key() + source = { + "name": name, + "key": key, + "created_at": 1000, + "enabled": True, + "revoked": False, + "revoked_at": None, + "stats": { + "segments_received": 0, + "entities_received": 0, + "facets_received": 0, + "imports_received": 0, + "config_received": 0, + }, + } + source.update(overrides) + return source + + +@pytest.fixture +def ingest_env(journal_env): + key = generate_key() + source = _source(key=key) + save_journal_source(source) + key_prefix = key[:8] + create_state_directory(journal_env, key_prefix) + + app = Flask(__name__) + app.config["TESTING"] = True + bp = Blueprint("import-test", __name__, url_prefix="/app/import") + register_ingest_routes(bp) + app.register_blueprint(bp) + + return { + "root": journal_env, + "key": key, + "key_prefix": key_prefix, + "source": source, + "client": app.test_client(), + } + + +def _post_config(client, key, key_prefix, config): + return client.post( + f"/app/import/journal/{key_prefix}/ingest/config", + headers={"Authorization": f"Bearer {key}"}, + json={"config": config}, + ) + + +def _sample_config(): + return { + "identity": {"name": "Remote User", "preferred": "Remote", "timezone": "UTC"}, + "convey": { + "password_hash": "secret_hash", + "secret": "secret_value", + "trust_localhost": True, + }, + "setup": {"completed_at": 12345}, + "env": {"API_KEY": "xyz"}, + "retention": {"days": 30}, + } + + +def _write_target_config(root, config): + config_dir = root / "config" + config_dir.mkdir(parents=True, exist_ok=True) + (config_dir / "journal.json").write_text( + json.dumps(config, ensure_ascii=False), encoding="utf-8" + ) + think.utils._journal_path_cache = None + + +def _read_json(path): + return json.loads(path.read_text(encoding="utf-8")) + + +def _read_log(key_prefix): + log_path = get_state_directory(key_prefix) / "config" / "log.jsonl" + if not log_path.exists(): + return [] + return [ + json.loads(line) + for line in log_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + + +def test_auth_required(ingest_env): + env = ingest_env + response = env["client"].post( + f"/app/import/journal/{env['key_prefix']}/ingest/config", + json={"config": {}}, + ) + assert response.status_code == 401 + + +def test_key_prefix_mismatch(ingest_env): + env = ingest_env + response = env["client"].post( + "/app/import/journal/deadbeef/ingest/config", + headers={"Authorization": f"Bearer {env['key']}"}, + json={"config": {}}, + ) + assert response.status_code == 403 + + +def test_invalid_json(ingest_env): + env = ingest_env + response = env["client"].post( + f"/app/import/journal/{env['key_prefix']}/ingest/config", + headers={"Authorization": f"Bearer {env['key']}"}, + data="not-json", + content_type="application/json", + ) + assert response.status_code == 400 + assert response.get_json() == {"error": "Invalid JSON body"} + + +def test_missing_config(ingest_env): + env = ingest_env + response = env["client"].post( + f"/app/import/journal/{env['key_prefix']}/ingest/config", + headers={"Authorization": f"Bearer {env['key']}"}, + json={}, + ) + assert response.status_code == 400 + assert response.get_json() == {"error": "Missing config object"} + + +def test_config_staged(ingest_env): + env = ingest_env + target_config = {"identity": {"name": "Local User"}, "retention": {"days": 90}} + _write_target_config(env["root"], target_config) + config = _sample_config() + + response = _post_config(env["client"], env["key"], env["key_prefix"], config) + body = response.get_json() + state_dir = get_state_directory(env["key_prefix"]) / "config" + source = load_journal_source(env["key"]) + + assert response.status_code == 200 + assert body == {"staged": True, "skipped": False, "diff_fields": 5} + assert (state_dir / "source_config.json").exists() + assert (state_dir / "diff.json").exists() + assert "last_hash" in _read_json(state_dir / "state.json") + assert _read_log(env["key_prefix"])[0]["action"] == "staged" + assert source["stats"]["config_received"] == 1 + + +def test_diff_categorization(ingest_env): + env = ingest_env + target_config = {"identity": {"name": "Local User"}, "retention": {"days": 90}} + _write_target_config(env["root"], target_config) + config = {"identity": {"name": "Remote User"}, "retention": {"days": 30}} + + response = _post_config(env["client"], env["key"], env["key_prefix"], config) + diff = _read_json(get_state_directory(env["key_prefix"]) / "config" / "diff.json") + + assert response.status_code == 200 + assert diff["identity.name"]["category"] == "transferable" + assert diff["retention.days"]["category"] == "preference" + + +def test_never_transfer_excluded(ingest_env): + env = ingest_env + _write_target_config(env["root"], {"identity": {"name": "Local User"}}) + config = _sample_config() + + response = _post_config(env["client"], env["key"], env["key_prefix"], config) + diff = _read_json(get_state_directory(env["key_prefix"]) / "config" / "diff.json") + + assert response.status_code == 200 + assert "convey.password_hash" not in diff + assert "convey.secret" not in diff + assert not any(key.startswith("env.") for key in diff) + + +def test_idempotent(ingest_env): + env = ingest_env + _write_target_config(env["root"], {"identity": {"name": "Local User"}}) + config = _sample_config() + + first = _post_config(env["client"], env["key"], env["key_prefix"], config) + second = _post_config(env["client"], env["key"], env["key_prefix"], config) + + assert first.status_code == 200 + assert second.status_code == 200 + assert second.get_json() == { + "staged": False, + "skipped": True, + "reason": "idempotent", + } + + +def test_config_always_staged(ingest_env): + env = ingest_env + _write_target_config(env["root"], {"identity": {"name": "Local User"}}) + config = {"identity": {"name": "Remote User"}} + + response = _post_config(env["client"], env["key"], env["key_prefix"], config) + + assert response.status_code == 200 + assert response.get_json()["staged"] is True diff --git a/tests/test_export.py b/tests/test_export.py index 9756e783c..4224f25b3 100644 --- a/tests/test_export.py +++ b/tests/test_export.py @@ -38,7 +38,9 @@ def _set_journal_override(monkeypatch, journal_path): think.utils._journal_path_cache = None -def _make_session(*, manifest_data=None, get_status=200, post_status=200, post_json=None): +def _make_session( + *, manifest_data=None, get_status=200, post_status=200, post_json=None +): mock = MagicMock(spec=requests.Session) mock.headers = {} @@ -140,6 +142,92 @@ def _facet_file_hash(path: Path) -> str: return hashlib.sha256(path.read_bytes()).hexdigest() +def _setup_imports(tmp_path): + """Create test import metadata dirs in journal fixture.""" + imports_dir = tmp_path / "imports" + imports_dir.mkdir(parents=True, exist_ok=True) + + dir1 = imports_dir / "20260101_090000" + dir1.mkdir() + import_json_1 = {"original_filename": "cal.zip", "file_size": 100} + imported_json_1 = { + "processed_timestamp": "20260101_090000", + "total_files_created": 1, + } + manifest_1 = [{"id": "event-0", "title": "Test Event"}] + (dir1 / "import.json").write_text(json.dumps(import_json_1), encoding="utf-8") + (dir1 / "imported.json").write_text(json.dumps(imported_json_1), encoding="utf-8") + (dir1 / "content_manifest.jsonl").write_text( + json.dumps(manifest_1[0]) + "\n", encoding="utf-8" + ) + + dir2 = imports_dir / "20260102_100000" + dir2.mkdir() + import_json_2 = {"original_filename": "chat.zip", "file_size": 200} + imported_json_2 = { + "processed_timestamp": "20260102_100000", + "total_files_created": 2, + } + manifest_2 = [{"id": "conv-0", "title": "Test Convo"}] + (dir2 / "import.json").write_text(json.dumps(import_json_2), encoding="utf-8") + (dir2 / "imported.json").write_text(json.dumps(imported_json_2), encoding="utf-8") + (dir2 / "content_manifest.jsonl").write_text( + json.dumps(manifest_2[0]) + "\n", encoding="utf-8" + ) + + (imports_dir / "plaud.json").write_text('{"last_sync": 123}', encoding="utf-8") + + source_dir = imports_dir / "abcd1234" + source_dir.mkdir() + (source_dir / "segments").mkdir() + + return { + "20260101_090000": { + "import_json": import_json_1, + "imported_json": imported_json_1, + "content_manifest": manifest_1, + }, + "20260102_100000": { + "import_json": import_json_2, + "imported_json": imported_json_2, + "content_manifest": manifest_2, + }, + } + + +def _import_hash(import_data): + """Compute hash matching export_imports algorithm.""" + hash_input = json.dumps( + { + "import_json": import_data["import_json"], + "imported_json": import_data["imported_json"], + "content_manifest": import_data["content_manifest"], + }, + sort_keys=True, + ensure_ascii=False, + ).encode() + return hashlib.sha256(hash_input).hexdigest() + + +def _setup_config(tmp_path): + """Create test config in journal fixture.""" + config_dir = tmp_path / "config" + config_dir.mkdir(parents=True, exist_ok=True) + config = { + "identity": {"name": "Test", "preferred": "Tester", "timezone": "UTC"}, + "convey": { + "password_hash": "secret_hash", + "secret": "secret_val", + "trust_localhost": True, + }, + "setup": {"completed_at": 12345}, + "env": {"KEY": "val"}, + "retention": {"days": 30}, + } + (config_dir / "journal.json").write_text(json.dumps(config), encoding="utf-8") + return config + + class TestExportSegments: def test_manifest_query_and_delta(self, tmp_path, monkeypatch): from observe.export import export_segments @@ -168,7 +256,9 @@ class TestExportSegments: mock_session = _make_session(manifest_data=manifest_data) with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=False) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=False + ) assert mock_session.post.call_count == 1 metadata = json.loads(mock_session.post.call_args.kwargs["data"]["metadata"]) @@ -183,7 +273,9 @@ class TestExportSegments: mock_session = _make_session(manifest_data={}) with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=True) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=True + ) assert mock_session.post.call_count == 0 output = capsys.readouterr().out @@ -217,7 +309,9 @@ class TestExportSegments: mock_session = _make_session(manifest_data=manifest_data) with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=True) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=True + ) assert mock_session.post.call_count == 0 output = capsys.readouterr().out @@ -251,7 +345,9 @@ class TestExportSegments: patch("observe.export.requests.Session", return_value=mock_session), patch("observe.export.time.sleep") as mock_sleep, ): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=False) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=False + ) assert mock_session.post.call_count == 3 assert mock_sleep.called @@ -282,7 +378,9 @@ class TestExportSegments: mock_session.get.side_effect = requests.ConnectionError with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=False) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=False + ) assert mock_session.post.call_count == 0 assert "Connection failed" in capsys.readouterr().out @@ -328,7 +426,9 @@ class TestExportSegments: mock_session = _make_session(manifest_data=manifest_data) with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=False) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=False + ) assert mock_session.post.call_count == 0 assert "up to date" in capsys.readouterr().out @@ -351,7 +451,9 @@ class TestExportSegments: mock_session.post.side_effect = [first, second] with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=False) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=False + ) assert mock_session.post.call_count == 2 output = capsys.readouterr().out @@ -367,7 +469,9 @@ class TestExportSegments: mock_session = _make_session(manifest_data={}) with patch("observe.export.requests.Session", return_value=mock_session): - export_segments("https://example.com", "test-key", ["20260413"], dry_run=False) + export_segments( + "https://example.com", "test-key", ["20260413"], dry_run=False + ) for call in mock_session.post.call_args_list: post_kwargs = call.kwargs @@ -572,7 +676,9 @@ class TestExportFacets: _set_journal_override(monkeypatch, tmp_path) post_json = {"created": 1, "merged": 0, "skipped": 0, "staged": 0, "errors": []} - mock_session = _make_session(manifest_data={"received": {}}, post_json=post_json) + mock_session = _make_session( + manifest_data={"received": {}}, post_json=post_json + ) with patch("observe.export.requests.Session", return_value=mock_session): export_facets("https://example.com", "test-key", dry_run=False) @@ -616,8 +722,14 @@ class TestExportFacets: { "name": "work", "files": [ - {"path": "entities/20260413.jsonl", "type": "detected_entities"}, - {"path": "entities/alice/entity.json", "type": "entity_relationship"}, + { + "path": "entities/20260413.jsonl", + "type": "detected_entities", + }, + { + "path": "entities/alice/entity.json", + "type": "entity_relationship", + }, { "path": "entities/alice/observations.jsonl", "type": "entity_observations", @@ -662,7 +774,10 @@ class TestExportFacets: assert metadata["facets"][0]["name"] == "work" assert metadata["facets"][0]["files"] == [ {"path": "entities/20260413.jsonl", "type": "detected_entities"}, - {"path": "entities/alice/observations.jsonl", "type": "entity_observations"}, + { + "path": "entities/alice/observations.jsonl", + "type": "entity_observations", + }, {"path": "todos/20260413.jsonl", "type": "todos"}, ] @@ -681,7 +796,9 @@ class TestExportFacets: output = capsys.readouterr().out assert "personal: 2 new, 0 changed, 0 unchanged" in output assert "work: 5 new, 0 changed, 0 unchanged" in output - assert "Dry run: 7 new files, 0 changed, 0 unchanged across 2 facet(s)" in output + assert ( + "Dry run: 7 new files, 0 changed, 0 unchanged across 2 facet(s)" in output + ) def test_idempotent(self, tmp_path, monkeypatch, capsys): from observe.export import export_facets @@ -696,7 +813,9 @@ class TestExportFacets: for file_path in sorted(facet_dir.rglob("*")): if file_path.is_file(): rel_path = file_path.relative_to(facet_dir).as_posix() - manifest_received[f"{facet_dir.name}/{rel_path}"] = _facet_file_hash(file_path) + manifest_received[f"{facet_dir.name}/{rel_path}"] = ( + _facet_file_hash(file_path) + ) mock_session = _make_session(manifest_data={"received": manifest_received}) with patch("observe.export.requests.Session", return_value=mock_session): @@ -765,7 +884,9 @@ class TestExportFacets: _set_journal_override(monkeypatch, tmp_path) post_json = {"created": 1, "merged": 0, "skipped": 0, "staged": 0, "errors": []} - mock_session = _make_session(manifest_data={"received": {}}, post_json=post_json) + mock_session = _make_session( + manifest_data={"received": {}}, post_json=post_json + ) with patch("observe.export.requests.Session", return_value=mock_session): export_facets("https://example.com", "test-key", dry_run=False) @@ -783,21 +904,29 @@ class TestExportFacets: _setup_facets(tmp_path) events_dir = tmp_path / "facets" / "work" / "events" events_dir.mkdir(parents=True) - (events_dir / "20260413.jsonl").write_text('{"event": "ignored"}\n', encoding="utf-8") + (events_dir / "20260413.jsonl").write_text( + '{"event": "ignored"}\n', encoding="utf-8" + ) _set_journal_override(monkeypatch, tmp_path) post_json = {"created": 1, "merged": 0, "skipped": 0, "staged": 0, "errors": []} - mock_session = _make_session(manifest_data={"received": {}}, post_json=post_json) + mock_session = _make_session( + manifest_data={"received": {}}, post_json=post_json + ) with patch("observe.export.requests.Session", return_value=mock_session): export_facets("https://example.com", "test-key", dry_run=False) calls_by_facet = { - json.loads(call.kwargs["data"]["metadata"])["facets"][0]["name"]: call.kwargs + json.loads(call.kwargs["data"]["metadata"])["facets"][0][ + "name" + ]: call.kwargs for call in mock_session.post.call_args_list } work_metadata = json.loads(calls_by_facet["work"]["data"]["metadata"]) - uploaded_paths = [entry["path"] for entry in work_metadata["facets"][0]["files"]] + uploaded_paths = [ + entry["path"] for entry in work_metadata["facets"][0]["files"] + ] assert "events/20260413.jsonl" not in uploaded_paths def test_response_errors_reported(self, tmp_path, monkeypatch, capsys): @@ -813,7 +942,9 @@ class TestExportFacets: "staged": 0, "errors": [{"facet": "work", "error": "entity merge conflict"}], } - mock_session = _make_session(manifest_data={"received": {}}, post_json=post_json) + mock_session = _make_session( + manifest_data={"received": {}}, post_json=post_json + ) with patch("observe.export.requests.Session", return_value=mock_session): export_facets("https://example.com", "test-key", dry_run=False) @@ -821,3 +952,176 @@ class TestExportFacets: output = capsys.readouterr().out assert "entity merge conflict" in output assert "error" in output.lower() + + +class TestExportImports: + def test_manifest_delta(self, tmp_path, monkeypatch): + from observe.export import export_imports + + imports = _setup_imports(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + manifest_data = { + "received": {"20260101_090000": _import_hash(imports["20260101_090000"])} + } + post_json = {"copied": 1, "staged": 0, "skipped": 0, "errors": []} + mock_session = _make_session(manifest_data=manifest_data, post_json=post_json) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_imports("https://example.com", "test-key", dry_run=False) + + assert mock_session.post.call_count == 1 + posted_data = mock_session.post.call_args.kwargs.get( + "json" + ) or mock_session.post.call_args[1].get("json") + posted_ids = [entry["id"] for entry in posted_data["imports"]] + assert posted_ids == ["20260102_100000"] + + def test_dry_run(self, tmp_path, monkeypatch, capsys): + from observe.export import export_imports + + _setup_imports(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + mock_session = _make_session(manifest_data={"received": {}}) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_imports("https://example.com", "test-key", dry_run=True) + + assert mock_session.post.call_count == 0 + output = capsys.readouterr().out + assert "2 new" in output + assert "0 changed" in output + + def test_idempotent(self, tmp_path, monkeypatch, capsys): + from observe.export import export_imports + + imports = _setup_imports(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + manifest_data = { + "received": { + import_id: _import_hash(import_data) + for import_id, import_data in imports.items() + } + } + mock_session = _make_session(manifest_data=manifest_data) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_imports("https://example.com", "test-key", dry_run=False) + + assert mock_session.post.call_count == 0 + assert "up to date" in capsys.readouterr().out + + def test_sync_state_excluded(self, tmp_path, monkeypatch): + from observe.export import export_imports + + _setup_imports(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + mock_session = _make_session( + manifest_data={"received": {}}, + post_json={"copied": 2, "staged": 0, "skipped": 0, "errors": []}, + ) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_imports("https://example.com", "test-key", dry_run=False) + + posted_data = mock_session.post.call_args.kwargs.get( + "json" + ) or mock_session.post.call_args[1].get("json") + posted_ids = {entry["id"] for entry in posted_data["imports"]} + assert "plaud.json" not in posted_ids + + def test_source_dir_excluded(self, tmp_path, monkeypatch): + from observe.export import export_imports + + _setup_imports(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + mock_session = _make_session( + manifest_data={"received": {}}, + post_json={"copied": 2, "staged": 0, "skipped": 0, "errors": []}, + ) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_imports("https://example.com", "test-key", dry_run=False) + + posted_data = mock_session.post.call_args.kwargs.get( + "json" + ) or mock_session.post.call_args[1].get("json") + posted_ids = {entry["id"] for entry in posted_data["imports"]} + assert "abcd1234" not in posted_ids + + +class TestExportConfig: + def test_config_export(self, tmp_path, monkeypatch): + from observe.export import export_config + + _setup_config(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + mock_session = _make_session( + manifest_data={}, + post_json={"staged": True, "skipped": False, "diff_fields": 3}, + ) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_config("https://example.com", "test-key", dry_run=False) + + assert mock_session.post.call_count == 1 + + def test_dry_run(self, tmp_path, monkeypatch, capsys): + from observe.export import export_config + + _setup_config(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + mock_session = _make_session(manifest_data={}) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_config("https://example.com", "test-key", dry_run=True) + + assert mock_session.post.call_count == 0 + assert "would send snapshot" in capsys.readouterr().out + + def test_idempotent(self, tmp_path, monkeypatch, capsys): + from observe.export import _strip_never_transfer, export_config + + config = _setup_config(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + stripped = _strip_never_transfer(config) + content_hash = hashlib.sha256( + json.dumps(stripped, sort_keys=True, ensure_ascii=False).encode() + ).hexdigest() + mock_session = _make_session(manifest_data={"last_hash": content_hash}) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_config("https://example.com", "test-key", dry_run=False) + + assert mock_session.post.call_count == 0 + assert "up to date" in capsys.readouterr().out + + def test_never_transfer_stripped(self, tmp_path, monkeypatch): + from observe.export import export_config + + _setup_config(tmp_path) + _set_journal_override(monkeypatch, tmp_path) + + mock_session = _make_session( + manifest_data={}, + post_json={"staged": True, "skipped": False, "diff_fields": 3}, + ) + + with patch("observe.export.requests.Session", return_value=mock_session): + export_config("https://example.com", "test-key", dry_run=False) + + posted_data = mock_session.post.call_args.kwargs.get( + "json" + ) or mock_session.post.call_args[1].get("json") + posted_config = posted_data["config"] + assert posted_config["convey"] == {"trust_localhost": True} + assert "setup" in posted_config + assert posted_config["setup"] == {} + assert "env" not in posted_config diff --git a/tests/test_imports_ingest.py b/tests/test_imports_ingest.py new file mode 100644 index 000000000..a7929df37 --- /dev/null +++ b/tests/test_imports_ingest.py @@ -0,0 +1,306 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import hashlib +import json +from importlib import import_module + +import pytest +from flask import Blueprint, Flask + +import convey.state + +journal_sources = import_module("apps.import.journal_sources") +ingest = import_module("apps.import.ingest") + +create_state_directory = journal_sources.create_state_directory +generate_key = journal_sources.generate_key +get_state_directory = journal_sources.get_state_directory +load_journal_source = journal_sources.load_journal_source +save_journal_source = journal_sources.save_journal_source +register_ingest_routes = ingest.register_ingest_routes + + +@pytest.fixture +def journal_env(tmp_path, monkeypatch): + monkeypatch.setattr(convey.state, "journal_root", str(tmp_path), raising=False) + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + (tmp_path / "apps" / "import" / "journal_sources").mkdir( + parents=True, exist_ok=True + ) + return tmp_path + + +def _source(name="test-source", key=None, **overrides): + if key is None: + key = generate_key() + source = { + "name": name, + "key": key, + "created_at": 1000, + "enabled": True, + "revoked": False, + "revoked_at": None, + "stats": { + "segments_received": 0, + "entities_received": 0, + "facets_received": 0, + "imports_received": 0, + "config_received": 0, + }, + } + source.update(overrides) + return source + + +@pytest.fixture +def ingest_env(journal_env): + key = generate_key() + source = _source(key=key) + save_journal_source(source) + key_prefix = key[:8] + create_state_directory(journal_env, key_prefix) + + app = Flask(__name__) + app.config["TESTING"] = True + bp = Blueprint("import-test", __name__, url_prefix="/app/import") + register_ingest_routes(bp) + app.register_blueprint(bp) + + return { + "root": journal_env, + "key": key, + "key_prefix": key_prefix, + "source": source, + "client": app.test_client(), + } + + +def _sample_import(import_id="20260101_090000"): + return { + "id": import_id, + "import_json": { + "original_filename": "test.zip", + "upload_timestamp": 1767258000000, + "upload_datetime": "2026-01-01T09:00:00", + "user_timestamp": import_id, + "file_size": 1234, + "mime_type": "application/zip", + "facet": "work", + "setting": "calendar", + "file_path": f"imports/{import_id}/test.zip", + }, + "imported_json": { + "processed_timestamp": import_id, + "processing_completed": "2026-01-01T09:10:00", + "total_files_created": 1, + "all_created_files": ["20260101/import.ics/090000_300/event.md"], + "segments": ["090000_300"], + "source_type": "ics", + "source_display": "Calendar", + "entries_written": 1, + "entities_seeded": 0, + "date_range": ["20260101", "20260101"], + "target_day": "20260101", + }, + "content_manifest": [ + { + "id": "event-0", + "title": "Test Event", + "date": "20260101", + "type": "event", + } + ], + } + + +def _import_hash(item: dict) -> str: + hash_input = json.dumps( + { + "import_json": item["import_json"], + "imported_json": item["imported_json"], + "content_manifest": item["content_manifest"], + }, + sort_keys=True, + ensure_ascii=False, + ).encode() + return hashlib.sha256(hash_input).hexdigest() + + +def _post_imports(client, key, key_prefix, imports_list): + return client.post( + f"/app/import/journal/{key_prefix}/ingest/imports", + headers={"Authorization": f"Bearer {key}"}, + json={"imports": imports_list}, + ) + + +def _read_state(key_prefix: str) -> dict: + state_path = get_state_directory(key_prefix) / "imports" / "state.json" + return json.loads(state_path.read_text(encoding="utf-8")) + + +def _read_log(key_prefix: str) -> list[dict]: + log_path = get_state_directory(key_prefix) / "imports" / "log.jsonl" + if not log_path.exists(): + return [] + return [ + json.loads(line) + for line in log_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + + +def test_auth_required(ingest_env): + env = ingest_env + response = env["client"].post( + f"/app/import/journal/{env['key_prefix']}/ingest/imports", + json={"imports": []}, + ) + assert response.status_code == 401 + + +def test_key_prefix_mismatch(ingest_env): + env = ingest_env + response = env["client"].post( + "/app/import/journal/deadbeef/ingest/imports", + headers={"Authorization": f"Bearer {env['key']}"}, + json={"imports": []}, + ) + assert response.status_code == 403 + + +def test_invalid_json(ingest_env): + env = ingest_env + response = env["client"].post( + f"/app/import/journal/{env['key_prefix']}/ingest/imports", + headers={"Authorization": f"Bearer {env['key']}"}, + data="not-json", + content_type="application/json", + ) + assert response.status_code == 400 + assert response.get_json() == {"error": "Invalid JSON body"} + + +def test_missing_imports_array(ingest_env): + env = ingest_env + response = env["client"].post( + f"/app/import/journal/{env['key_prefix']}/ingest/imports", + headers={"Authorization": f"Bearer {env['key']}"}, + json={}, + ) + assert response.status_code == 400 + assert response.get_json() == {"error": "Missing imports array"} + + +def test_copy_new_import(ingest_env): + env = ingest_env + item = _sample_import() + response = _post_imports(env["client"], env["key"], env["key_prefix"], [item]) + body = response.get_json() + import_dir = env["root"] / "imports" / item["id"] + + assert response.status_code == 200 + assert body == {"copied": 1, "skipped": 0, "staged": 0, "errors": []} + assert ( + json.loads((import_dir / "import.json").read_text(encoding="utf-8")) + == item["import_json"] + ) + assert ( + json.loads((import_dir / "imported.json").read_text(encoding="utf-8")) + == item["imported_json"] + ) + assert [ + json.loads(line) + for line in (import_dir / "content_manifest.jsonl") + .read_text(encoding="utf-8") + .splitlines() + if line.strip() + ] == item["content_manifest"] + assert _read_state(env["key_prefix"]) == { + "received": {item["id"]: _import_hash(item)} + } + assert _read_log(env["key_prefix"])[0]["action"] == "copied" + source = load_journal_source(env["key"]) + assert source["stats"]["imports_received"] == 1 + + +def test_dedup_identical(ingest_env): + env = ingest_env + item = _sample_import() + + first = _post_imports(env["client"], env["key"], env["key_prefix"], [item]) + second = _post_imports(env["client"], env["key"], env["key_prefix"], [item]) + + assert first.status_code == 200 + assert second.status_code == 200 + assert second.get_json() == {"copied": 0, "skipped": 1, "staged": 0, "errors": []} + assert _read_log(env["key_prefix"])[-1]["reason"] == "idempotent" + + +def test_id_collision(ingest_env): + env = ingest_env + item = _sample_import() + (env["root"] / "imports" / item["id"]).mkdir(parents=True, exist_ok=True) + + response = _post_imports(env["client"], env["key"], env["key_prefix"], [item]) + body = response.get_json() + staged_path = ( + get_state_directory(env["key_prefix"]) + / "imports" + / "staged" + / f"{item['id']}.json" + ) + + assert response.status_code == 200 + assert body == {"copied": 0, "skipped": 0, "staged": 1, "errors": []} + assert staged_path.exists() + assert ( + json.loads(staged_path.read_text(encoding="utf-8"))["reason"] == "id_collision" + ) + + +def test_multiple_imports(ingest_env): + env = ingest_env + first = _sample_import("20260101_090000") + second = _sample_import("20260101_100000") + third = _sample_import("20260101_110000") + _post_imports(env["client"], env["key"], env["key_prefix"], [first]) + (env["root"] / "imports" / third["id"]).mkdir(parents=True, exist_ok=True) + + response = _post_imports( + env["client"], env["key"], env["key_prefix"], [first, second, third] + ) + + assert response.status_code == 200 + assert response.get_json() == {"copied": 1, "skipped": 1, "staged": 1, "errors": []} + + +def test_state_manifest(ingest_env): + env = ingest_env + item = _sample_import() + + response = _post_imports(env["client"], env["key"], env["key_prefix"], [item]) + + assert response.status_code == 200 + assert _read_state(env["key_prefix"]) == { + "received": {item["id"]: _import_hash(item)} + } + + +def test_stats_update(ingest_env): + env = ingest_env + first = _sample_import("20260101_090000") + second = _sample_import("20260101_100000") + collision = _sample_import("20260101_110000") + (env["root"] / "imports" / collision["id"]).mkdir(parents=True, exist_ok=True) + + response = _post_imports( + env["client"], env["key"], env["key_prefix"], [first, second, collision] + ) + source = load_journal_source(env["key"]) + + assert response.status_code == 200 + assert source["stats"]["imports_received"] == 2