diff --git a/apps/observer/events.py b/apps/observer/events.py index ea2ec089b..437fccc56 100644 --- a/apps/observer/events.py +++ b/apps/observer/events.py @@ -57,3 +57,51 @@ def handle_observed(ctx: EventContext) -> None: logger.debug( f"Recorded observed status for observer {observer_name}: {day}/{segment}" ) + + +@on_event("observe", "transferred") +def handle_transferred(ctx: EventContext) -> None: + """Handle observe.transferred events for transfer-originated segments. + + When a transferred segment is received, append a 'transferred' record + to the observer's sync history, increment stats, and queue an indexer + rescan to pick up the new content. + """ + observer_name = ctx.msg.get("observer") + if not observer_name: + return + + segment = ctx.msg.get("segment") + day = ctx.msg.get("day") + if not segment or not day: + logger.warning( + f"observe.transferred missing segment/day for observer {observer_name}" + ) + return + + observer = find_observer_by_name(observer_name) + if not observer: + logger.debug(f"Observer not found for transferred event: {observer_name}") + return + + key_prefix = observer.get("key", "")[:8] + if not key_prefix: + return + + record = { + "ts": now_ms(), + "type": "transferred", + "segment": segment, + } + append_history_record(key_prefix, day, record) + + increment_stat(key_prefix, "segments_transferred") + + # Queue indexer rescan to pick up transferred content + from think.callosum import callosum_send + + callosum_send("supervisor", "request", cmd=["sol", "indexer", "--rescan"]) + + logger.debug( + f"Recorded transferred status for observer {observer_name}: {day}/{segment}" + ) diff --git a/apps/observer/routes.py b/apps/observer/routes.py index 24765fbb9..ffe83f991 100644 --- a/apps/observer/routes.py +++ b/apps/observer/routes.py @@ -6,6 +6,8 @@ Provides endpoints for: - Managing observer registrations (UI) - Receiving file uploads from observers (ingest) +- Receiving transferred segments from other instances (transfer ingest) +- Serving segment manifests for transfer diffing - Relaying events from observers to local Callosum - Retrieving segment upload history for sync verification """ @@ -15,6 +17,7 @@ from __future__ import annotations import base64 import json import logging +import platform import re import secrets from pathlib import Path @@ -28,14 +31,16 @@ from convey import emit from observe.utils import ( MAX_SEGMENT_ATTEMPTS, compute_bytes_sha256, + compute_file_sha256, find_available_segment, ) from think.streams import stream_name, update_stream, write_segment_stream -from think.utils import day_path, now_ms, segment_path +from think.utils import day_path, iter_segments, now_ms, segment_path from .utils import ( append_history_record, find_segment_by_sha256, + get_hist_dir, get_observers_dir, list_observers, load_history, @@ -301,96 +306,44 @@ def _save_to_failed( # === Ingest API (key-protected) === -@observer_bp.route("/ingest", methods=["POST"]) -@observer_bp.route("/ingest/", methods=["POST"]) -def ingest_upload(key: str | None = None) -> Any: - """Receive file uploads from observer. - - Expects multipart form with: - - segment: Segment key (HHMMSS_LEN) - - day: Day string (YYYYMMDD) - - files: One or more media files - - host: (optional) Hostname of observer - - platform: (optional) Platform of observer - - meta: (optional) JSON-encoded metadata dict (facet, setting, etc.) - - Writes files to journal and emits observe.observing event. - Host/platform are merged into meta (meta values take precedence). - - Returns status: - - "ok": New segment accepted - - "duplicate": All files already received (no processing triggered) - - "collision": New segment saved with adjusted key (directory conflict) +def _process_ingest_files( + observer: dict, + key_prefix: str, + segment: str, + day: str, + stream: str, + uploaded_files, + *, + source: str | None = None, +) -> tuple[dict, int]: + """Shared ingest pipeline: read/hash files, dedup, deconflict, save, record history, update stats. + + Parameters + ---------- + observer : dict + Observer metadata dict (must include 'stats', 'name', 'last_seen', etc.) + key_prefix : str + First 8 chars of observer key. + segment : str + Requested segment key (HHMMSS_LEN format). + day : str + Day string (YYYYMMDD format). + stream : str + Stream name (already resolved by caller). + uploaded_files : list + List of Flask FileStorage objects from request.files.getlist("files"). + source : str or None + If provided, added as "source" field to history record (e.g., "transfer"). + + Returns + ------- + tuple of (dict, int) + Response body dict and HTTP status code. """ - # Extract key from Bearer header (primary) or URL path (legacy) - auth_key = _get_key(key) - if not auth_key: - return jsonify({"error": "Authorization required"}), 401 - - # Validate key - observer = load_observer(auth_key) - if not observer: - return jsonify({"error": "Invalid key"}), 401 - - if observer.get("revoked", False): - return jsonify({"error": "Observer revoked"}), 403 - - if not observer.get("enabled", True): - return jsonify({"error": "Observer disabled"}), 403 - - # Get segment, day, and host info from form - segment = request.form.get("segment", "").strip() - day = request.form.get("day", "").strip() - host = request.form.get("host", "").strip() - platform = request.form.get("platform", "").strip() - meta_str = request.form.get("meta", "").strip() - - # Parse meta JSON and merge host/platform (meta values take precedence) - meta: dict = {} - if meta_str: - try: - meta = json.loads(meta_str) - except json.JSONDecodeError: - logger.warning(f"Invalid meta JSON from observer: {meta_str[:100]}") - if host and "host" not in meta: - meta["host"] = host - if platform and "platform" not in meta: - meta["platform"] = platform - - # Warn if client hostname differs from registered observer name - effective_host = meta.get("host", host) - observer_name = observer.get("name", "") - if effective_host and effective_host != observer_name: - logger.warning( - f"Observer '{observer_name}' ({auth_key[:8]}) connecting from host " - f"'{effective_host}' — hostname differs from registered name. " - f"Use `sol observer rename` to update if the host was renamed." - ) - - if not segment: - return jsonify({"error": "Missing segment"}), 400 - if not day: - return jsonify({"error": "Missing day"}), 400 - - # Validate segment format (HHMMSS_LEN) - if not re.match(r"^\d{6}_\d+$", segment): - return jsonify({"error": "Invalid segment format"}), 400 - - # Validate day format (YYYYMMDD) - if not re.match(r"^\d{8}$", day): - return jsonify({"error": "Invalid day format"}), 400 - - # Get uploaded files - files = request.files.getlist("files") - if not files: - return jsonify({"error": "No files uploaded"}), 400 - - key_prefix = auth_key[:8] - # Read file contents into memory and compute SHA256 before saving # This allows duplicate detection without writing to disk file_data = [] # List of (submitted_filename, simple_filename, content, sha256) - for upload in files: + for upload in uploaded_files: if not upload.filename: continue @@ -411,7 +364,7 @@ def ingest_upload(key: str | None = None) -> Any: file_data.append((submitted_filename, simple_filename, content, sha256)) if not file_data: - return jsonify({"error": "No valid files uploaded"}), 400 + return {"error": "No valid files uploaded"}, 400 # Check for duplicate submission by SHA256 incoming_sha256s = {fd[3] for fd in file_data} @@ -420,46 +373,32 @@ def ingest_upload(key: str | None = None) -> Any: ) if existing_segment: - # Full duplicate - all files already exist in an existing segment logger.info( f"Duplicate segment rejected: {day}/{segment} from {observer.get('name')} " f"(matches existing {existing_segment})" ) - # Update last_seen and increment duplicates_rejected stat observer["last_seen"] = now_ms() observer["stats"]["duplicates_rejected"] = ( observer["stats"].get("duplicates_rejected", 0) + 1 ) save_observer(observer) - return jsonify( + return ( { "status": "duplicate", "existing_segment": existing_segment, "message": "All files already received", - } + }, + 200, ) - # Log partial match context if some files already exist partial_match = bool(matched_sha256s) # Ensure day directory exists day_dir = day_path(day) day_dir.mkdir(parents=True, exist_ok=True) - # Determine stream name: trust client-provided stream in meta if valid, - # otherwise derive from observer registration name. - # Deriving from observer name via stream_name(observer=...) calls _strip_hostname, - # which strips qualifiers like ".tmux" — so "fedora.tmux" becomes "fedora", - # colliding both observers into one stream. - client_stream = meta.get("stream", "").strip() - observer_name = observer.get("name", "unknown") - if client_stream and re.match(r"^[a-z0-9][a-z0-9._-]*$", client_stream): - stream = client_stream - else: - stream = stream_name(observer=observer_name) - # Find available segment key within the stream directory stream_dir = day_dir / stream stream_dir.mkdir(parents=True, exist_ok=True) @@ -468,28 +407,25 @@ def ingest_upload(key: str | None = None) -> Any: available_segment = find_available_segment(stream_dir, segment) if available_segment is None: - # Exhausted attempts, save to failed directory logger.error( f"No available segment slot for {day}/{stream}/{segment} from " - f"{observer_name} after {MAX_SEGMENT_ATTEMPTS} attempts" + f"{observer.get('name', 'unknown')} after {MAX_SEGMENT_ATTEMPTS} attempts" ) failed_dir = _save_to_failed(day_dir, file_data, segment) return ( - jsonify( - { - "status": "failed", - "error": f"No available segment slot after {MAX_SEGMENT_ATTEMPTS} attempts", - "failed_path": str(failed_dir.relative_to(day_dir.parent)), - } - ), + { + "status": "failed", + "error": f"No available segment slot after {MAX_SEGMENT_ATTEMPTS} attempts", + "failed_path": str(failed_dir.relative_to(day_dir.parent)), + }, 507, - ) # Insufficient Storage + ) segment = available_segment if segment != original_segment: logger.info( f"Segment collision resolved: {original_segment} -> {segment} " - f"for observer {observer_name}" + f"for observer {observer.get('name', 'unknown')}" ) # Create segment directory for files (under stream) @@ -526,12 +462,11 @@ def ingest_upload(key: str | None = None) -> Any: logger.info(f"Saved {simple_filename} to {segment_dir}") except OSError as e: logger.error(f"Failed to save {simple_filename}: {e}") - return jsonify({"error": f"Failed to save {simple_filename}"}), 500 + return {"error": f"Failed to save {simple_filename}"}, 500 if not saved_files: - return jsonify({"error": "No valid files saved"}), 400 + return {"error": "No valid files saved"}, 400 - # Write sync history record sync_record = { "ts": now_ms(), "segment": segment, @@ -541,11 +476,11 @@ def ingest_upload(key: str | None = None) -> Any: if segment != original_segment: sync_record["segment_original"] = original_segment if partial_match: - # Log which SHA256s matched existing files (for debugging/audit) sync_record["partial_match_sha256s"] = list(matched_sha256s) + if source: + sync_record["source"] = source append_history_record(key_prefix, day, sync_record) - # Update observer stats observer["last_seen"] = now_ms() observer["last_segment"] = segment observer["stats"]["segments_received"] = ( @@ -556,6 +491,123 @@ def ingest_upload(key: str | None = None) -> Any: ) save_observer(observer) + status = "collision" if segment != original_segment else "ok" + return { + "status": status, + "segment": segment, + "files": saved_files, + "bytes": total_bytes, + }, 200 + + +@observer_bp.route("/ingest", methods=["POST"]) +@observer_bp.route("/ingest/", methods=["POST"]) +def ingest_upload(key: str | None = None) -> Any: + """Receive file uploads from observer. + + Expects multipart form with: + - segment: Segment key (HHMMSS_LEN) + - day: Day string (YYYYMMDD) + - files: One or more media files + - host: (optional) Hostname of observer + - platform: (optional) Platform of observer + - meta: (optional) JSON-encoded metadata dict (facet, setting, etc.) + + Writes files to journal and emits observe.observing event. + Host/platform are merged into meta (meta values take precedence). + + Returns status: + - "ok": New segment accepted + - "duplicate": All files already received (no processing triggered) + - "collision": New segment saved with adjusted key (directory conflict) + """ + # Extract key from Bearer header (primary) or URL path (legacy) + auth_key = _get_key(key) + if not auth_key: + return jsonify({"error": "Authorization required"}), 401 + + # Validate key + observer = load_observer(auth_key) + if not observer: + return jsonify({"error": "Invalid key"}), 401 + + if observer.get("revoked", False): + return jsonify({"error": "Observer revoked"}), 403 + + if not observer.get("enabled", True): + return jsonify({"error": "Observer disabled"}), 403 + + # Get segment, day, and host info from form + segment = request.form.get("segment", "").strip() + day = request.form.get("day", "").strip() + host = request.form.get("host", "").strip() + platform = request.form.get("platform", "").strip() + meta_str = request.form.get("meta", "").strip() + + # Parse meta JSON and merge host/platform (meta values take precedence) + meta: dict = {} + if meta_str: + try: + meta = json.loads(meta_str) + except json.JSONDecodeError: + logger.warning(f"Invalid meta JSON from observer: {meta_str[:100]}") + if host and "host" not in meta: + meta["host"] = host + if platform and "platform" not in meta: + meta["platform"] = platform + + # Warn if client hostname differs from registered observer name + effective_host = meta.get("host", host) + observer_name = observer.get("name", "") + if effective_host and effective_host != observer_name: + logger.warning( + f"Observer '{observer_name}' ({auth_key[:8]}) connecting from host " + f"'{effective_host}' — hostname differs from registered name. " + f"Use `sol observer rename` to update if the host was renamed." + ) + + if not segment: + return jsonify({"error": "Missing segment"}), 400 + if not day: + return jsonify({"error": "Missing day"}), 400 + + # Validate segment format (HHMMSS_LEN) + if not re.match(r"^\d{6}_\d+$", segment): + return jsonify({"error": "Invalid segment format"}), 400 + + # Validate day format (YYYYMMDD) + if not re.match(r"^\d{8}$", day): + return jsonify({"error": "Invalid day format"}), 400 + + # Get uploaded files + files = request.files.getlist("files") + if not files: + return jsonify({"error": "No files uploaded"}), 400 + + key_prefix = auth_key[:8] + + # Determine stream name: trust client-provided stream in meta if valid, + # otherwise derive from observer registration name. + # Deriving from observer name via stream_name(observer=...) calls _strip_hostname, + # which strips qualifiers like ".tmux" — so "fedora.tmux" becomes "fedora", + # colliding both observers into one stream. + client_stream = meta.get("stream", "").strip() + observer_name = observer.get("name", "unknown") + if client_stream and re.match(r"^[a-z0-9][a-z0-9._-]*$", client_stream): + stream = client_stream + else: + stream = stream_name(observer=observer_name) + + body, status = _process_ingest_files( + observer, key_prefix, segment, day, stream, files + ) + if status != 200 or body.get("status") == "duplicate": + return jsonify(body), status + + segment = body["segment"] + saved_files = body["files"] + segment_dir = segment_path(day, segment, stream) + # Write stream identity for this segment try: result = update_stream(stream, day, segment, type="observer") @@ -588,21 +640,167 @@ def ingest_upload(key: str | None = None) -> Any: logger.info( f"Received {len(saved_files)} files for {day}/{segment} from {observer.get('name')}" ) + return jsonify(body), status - # Determine response status - if segment != original_segment: - status = "collision" - else: - status = "ok" - return jsonify( - { - "status": status, - "segment": segment, - "files": saved_files, - "bytes": total_bytes, - } +@observer_bp.route("/ingest//transfer", methods=["POST"]) +def ingest_transfer(key: str) -> Any: + """Receive transferred file uploads from another solstone instance.""" + auth_key = _get_key(key) + if not auth_key: + return jsonify({"error": "Authorization required"}), 401 + + observer = load_observer(auth_key) + if not observer: + return jsonify({"error": "Invalid key"}), 401 + + if observer.get("revoked", False): + return jsonify({"error": "Observer revoked"}), 403 + + if not observer.get("enabled", True): + return jsonify({"error": "Observer disabled"}), 403 + + segment = request.form.get("segment", "").strip() + day = request.form.get("day", "").strip() + stream = request.form.get("stream", "").strip() + host = request.form.get("host", "").strip() + platform_name = request.form.get("platform", "").strip() + meta_str = request.form.get("meta", "").strip() + + meta: dict = {} + if meta_str: + try: + meta = json.loads(meta_str) + except json.JSONDecodeError: + logger.warning(f"Invalid meta JSON from observer: {meta_str[:100]}") + if host and "host" not in meta: + meta["host"] = host + if platform_name and "platform" not in meta: + meta["platform"] = platform_name + + if not segment: + return jsonify({"error": "Missing segment"}), 400 + if not day: + return jsonify({"error": "Missing day"}), 400 + if not stream: + return jsonify({"error": "Missing stream"}), 400 + if not re.match(r"^\d{6}_\d+$", segment): + return jsonify({"error": "Invalid segment format"}), 400 + if not re.match(r"^\d{8}$", day): + return jsonify({"error": "Invalid day format"}), 400 + if not re.match(r"^[a-z0-9][a-z0-9._-]*$", stream): + return jsonify({"error": "Invalid stream format"}), 400 + + files = request.files.getlist("files") + if not files: + return jsonify({"error": "No files uploaded"}), 400 + + key_prefix = auth_key[:8] + body, status = _process_ingest_files( + observer, + key_prefix, + segment, + day, + stream, + files, + source="transfer", ) + if status != 200 or body.get("status") == "duplicate": + return jsonify(body), status + + observer_name = observer.get("name", "") + event_fields: dict[str, Any] = { + "segment": body["segment"], + "day": day, + "files": body["files"], + "observer": observer_name, + "stream": stream, + } + if meta: + event_fields["meta"] = meta + emit("observe", "transferred", **event_fields) + + return jsonify(body), status + + +@observer_bp.route("/ingest//manifest", methods=["GET"]) +def ingest_manifest(key: str) -> Any: + """List available manifest days for an observer.""" + auth_key = _get_key(key) + if not auth_key: + return jsonify({"error": "Authorization required"}), 401 + + observer = load_observer(auth_key) + if not observer: + return jsonify({"error": "Invalid key"}), 401 + + if observer.get("revoked", False): + return jsonify({"error": "Observer revoked"}), 403 + + if not observer.get("enabled", True): + return jsonify({"error": "Observer disabled"}), 403 + + key_prefix = auth_key[:8] + hist_dir = get_hist_dir(key_prefix, ensure_exists=False) + if not hist_dir.exists(): + return jsonify({"days": {}}) + + days: dict[str, dict[str, int]] = {} + for hist_path in sorted(hist_dir.glob("*.jsonl")): + records = load_history(key_prefix, hist_path.stem) + segments = { + record.get("segment", "") + for record in records + if not record.get("type") and record.get("segment") + } + days[hist_path.stem] = {"segments": len(segments)} + + return jsonify({"days": days}) + + +@observer_bp.route("/ingest//manifest/", methods=["GET"]) +def ingest_manifest_day(key: str, day: str) -> Any: + """Return a transfer manifest for all segments on a given day.""" + auth_key = _get_key(key) + if not auth_key: + return jsonify({"error": "Authorization required"}), 401 + + observer = load_observer(auth_key) + if not observer: + return jsonify({"error": "Invalid key"}), 401 + + if observer.get("revoked", False): + return jsonify({"error": "Observer revoked"}), 403 + + if not observer.get("enabled", True): + return jsonify({"error": "Observer disabled"}), 403 + + if not re.match(r"^\d{8}$", day): + return jsonify({"error": "Invalid day format"}), 400 + + manifest = { + "version": 1, + "day": day, + "created_at": now_ms(), + "host": platform.node() or "unknown", + "segments": {}, + } + + for stream, seg_key, seg_path in iter_segments(day): + arc_key = f"{stream}/{seg_key}" + files = [] + for file_path in sorted(seg_path.iterdir()): + if file_path.is_file(): + files.append( + { + "name": file_path.name, + "sha256": compute_file_sha256(file_path), + "size": file_path.stat().st_size, + } + ) + manifest["segments"][arc_key] = {"files": files} + + return jsonify(manifest) @observer_bp.route("/ingest/event", methods=["POST"]) diff --git a/apps/observer/tests/test_events.py b/apps/observer/tests/test_events.py index 753397fb8..92dcbf5a3 100644 --- a/apps/observer/tests/test_events.py +++ b/apps/observer/tests/test_events.py @@ -10,7 +10,7 @@ import json import pytest from apps.events import EventContext -from apps.observer.events import handle_observed +from apps.observer.events import handle_observed, handle_transferred @pytest.fixture @@ -75,7 +75,9 @@ class TestHandleObserved: handle_observed(ctx) # Check history was written - hist_path = observer_journal.observers_dir / "testkey1" / "hist" / "20250103.jsonl" + hist_path = ( + observer_journal.observers_dir / "testkey1" / "hist" / "20250103.jsonl" + ) assert hist_path.exists() with open(hist_path) as f: @@ -108,7 +110,9 @@ class TestHandleObserved: handle_observed(ctx) # Check all records written - hist_path = observer_journal.observers_dir / "testkey1" / "hist" / "20250103.jsonl" + hist_path = ( + observer_journal.observers_dir / "testkey1" / "hist" / "20250103.jsonl" + ) with open(hist_path) as f: lines = f.readlines() @@ -204,3 +208,50 @@ class TestHandleObserved: # No history should be created hist_dir = observer_journal.observers_dir / "testkey1" / "hist" assert not hist_dir.exists() + + def test_handle_transferred(self, observer_journal, monkeypatch): + """Handler records transferred status, stats, and queues rescan.""" + import think.callosum as callosum_module + + calls = [] + monkeypatch.setattr( + callosum_module, + "callosum_send", + lambda *a, **kw: calls.append((a, kw)) or True, + ) + + ctx = EventContext( + msg={ + "tract": "observe", + "event": "transferred", + "observer": "test-observer", + "segment": "120000_300", + "day": "20250103", + }, + app="observer", + tract="observe", + event="transferred", + ) + + handle_transferred(ctx) + + hist_path = ( + observer_journal.observers_dir / "testkey1" / "hist" / "20250103.jsonl" + ) + assert hist_path.exists() + with open(hist_path) as f: + record = json.loads(f.readline()) + + assert record["type"] == "transferred" + assert record["segment"] == "120000_300" + + with open(observer_journal.observer_path) as f: + data = json.load(f) + assert data["stats"]["segments_transferred"] == 1 + + assert calls == [ + ( + ("supervisor", "request"), + {"cmd": ["sol", "indexer", "--rescan"]}, + ) + ] diff --git a/apps/observer/tests/test_routes.py b/apps/observer/tests/test_routes.py index fc3a65ccf..eaf908ad6 100644 --- a/apps/observer/tests/test_routes.py +++ b/apps/observer/tests/test_routes.py @@ -1255,7 +1255,9 @@ def test_api_list_includes_segments_observed_stat(observer_env): assert "segments_observed" not in data[0]["stats"] # Manually add segments_observed stat - observer_path = env.journal / "apps" / "observer" / "observers" / f"{key_prefix}.json" + observer_path = ( + env.journal / "apps" / "observer" / "observers" / f"{key_prefix}.json" + ) with open(observer_path) as f: observer_data = json.load(f) observer_data["stats"]["segments_observed"] = 5 @@ -1690,3 +1692,362 @@ def test_ingest_stream_qualifier_preserved(observer_env): assert not ( env.journal / "20250103" / "fedora" / "120000_300" / "tmux.jsonl" ).exists() + + +def test_transfer_success(observer_env): + """Test successful transfer upload.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + test_data = b"transferred audio content" + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(test_data), "audio.flac"), + }, + ) + assert resp.status_code == 200 + data = resp.get_json() + assert data["status"] == "ok" + assert data["segment"] == "120000_300" + assert data["files"] == ["audio.flac"] + assert data["bytes"] == len(test_data) + + expected_file = ( + env.journal / "20250103" / "remote.host" / "120000_300" / "audio.flac" + ) + assert expected_file.exists() + assert expected_file.read_bytes() == test_data + + +def test_transfer_requires_stream(observer_env): + """Test that transfer requires stream.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-stream-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "files": (io.BytesIO(b"content"), "audio.flac"), + }, + ) + assert resp.status_code == 400 + assert resp.get_json()["error"] == "Missing stream" + + +def test_transfer_invalid_stream(observer_env): + """Test that transfer validates stream format.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-invalid-stream"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "INVALID!", + "files": (io.BytesIO(b"content"), "audio.flac"), + }, + ) + assert resp.status_code == 400 + assert resp.get_json()["error"] == "Invalid stream format" + + +def test_transfer_duplicate_detection(observer_env): + """Test transfer duplicate detection.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-duplicate-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + test_data = b"duplicate transfer content" + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(test_data), "audio.flac"), + }, + ) + assert resp.status_code == 200 + assert resp.get_json()["status"] == "ok" + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(test_data), "audio.flac"), + }, + ) + assert resp.status_code == 200 + data = resp.get_json() + assert data["status"] == "duplicate" + assert data["existing_segment"] == "120000_300" + + +def test_transfer_deconfliction(observer_env): + """Test transfer deconflicts existing segment directories.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-collision-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + stream_dir = env.journal / "20250103" / "remote.host" + stream_dir.mkdir(parents=True) + (stream_dir / "120000_300").mkdir() + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(b"collision content"), "audio.flac"), + }, + ) + assert resp.status_code == 200 + data = resp.get_json() + assert data["status"] == "collision" + assert data["segment"] != "120000_300" + assert (stream_dir / data["segment"] / "audio.flac").exists() + + +def test_transfer_emits_transferred_event(observer_env, monkeypatch): + """Test transfer emits observe.transferred.""" + env = observer_env() + + import apps.observer.routes as routes_module + + calls = [] + + def mock_emit(*args, **kwargs): + calls.append((args, kwargs)) + + monkeypatch.setattr(routes_module, "emit", mock_emit) + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-event-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(b"event content"), "audio.flac"), + }, + ) + assert resp.status_code == 200 + assert len(calls) == 1 + assert calls[0][0] == ("observe", "transferred") + + +def test_transfer_does_not_emit_observing(observer_env, monkeypatch): + """Test transfer does not emit observe.observing.""" + env = observer_env() + + import apps.observer.routes as routes_module + + calls = [] + + def mock_emit(*args, **kwargs): + calls.append((args, kwargs)) + + monkeypatch.setattr(routes_module, "emit", mock_emit) + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-no-observing-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(b"event content"), "audio.flac"), + }, + ) + assert resp.status_code == 200 + assert all(args[1] != "observing" for args, _kwargs in calls) + + +def test_transfer_history_record(observer_env): + """Test transfer upload history records source='transfer'.""" + from apps.observer.utils import load_history + + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "transfer-history-test"}, + content_type="application/json", + ) + data = resp.get_json() + key = data["key"] + key_prefix = data["key_prefix"] + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(b"history content"), "audio.flac"), + }, + ) + assert resp.status_code == 200 + + records = load_history(key_prefix, "20250103") + upload_record = next(record for record in records if not record.get("type")) + assert upload_record["source"] == "transfer" + + +def test_transfer_auth_required(observer_env): + """Test transfer rejects invalid path key without auth header.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/ingest/badkey/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(b"content"), "audio.flac"), + }, + ) + assert resp.status_code == 401 + + +def test_transfer_invalid_key(observer_env): + """Test transfer rejects invalid key.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/ingest/not-a-real-key/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": (io.BytesIO(b"content"), "audio.flac"), + }, + ) + assert resp.status_code == 401 + + +def test_manifest_day_listing(observer_env): + """Test manifest day listing from observer history.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "manifest-list-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + resp = env.client.post( + "/app/observer/ingest", + headers={"Authorization": f"Bearer {key}"}, + data={ + "day": "20250103", + "segment": "120000_300", + "files": (io.BytesIO(b"manifest content"), "audio.flac"), + }, + ) + assert resp.status_code == 200 + + resp = env.client.get(f"/app/observer/ingest/{key}/manifest") + assert resp.status_code == 200 + assert resp.get_json() == {"days": {"20250103": {"segments": 1}}} + + +def test_manifest_per_day(observer_env): + """Test per-day manifest format matches transfer manifest v1.""" + env = observer_env() + + resp = env.client.post( + "/app/observer/api/create", + json={"name": "manifest-day-test"}, + content_type="application/json", + ) + key = resp.get_json()["key"] + + resp = env.client.post( + f"/app/observer/ingest/{key}/transfer", + data={ + "day": "20250103", + "segment": "120000_300", + "stream": "remote.host", + "files": [ + (io.BytesIO(b"audio bytes"), "audio.flac"), + (io.BytesIO(b"screen bytes"), "screen.webm"), + ], + }, + ) + assert resp.status_code == 200 + + resp = env.client.get(f"/app/observer/ingest/{key}/manifest/20250103") + assert resp.status_code == 200 + data = resp.get_json() + + assert data["version"] == 1 + assert data["day"] == "20250103" + assert isinstance(data["created_at"], int) + assert "host" in data + assert "remote.host/120000_300" in data["segments"] + + files = data["segments"]["remote.host/120000_300"]["files"] + assert len(files) == 2 + for file_info in files: + assert set(file_info) == {"name", "sha256", "size"} + assert len(file_info["sha256"]) == 64 + + +def test_manifest_auth_required(observer_env): + """Test manifest endpoint rejects invalid key.""" + env = observer_env() + + resp = env.client.get("/app/observer/ingest/badkey/manifest") + assert resp.status_code == 401