From d46f5947d8cb4bcddaa4f4b29fcb3f5d15a268a1 Mon Sep 17 00:00:00 2001 From: Jer Miller Date: Sun, 14 Jun 2026 17:13:34 -0600 Subject: [PATCH] =?UTF-8?q?refactor(observer):=20retire=20legacy=20DL=20ke?= =?UTF-8?q?y-in-URL=20ingest=20surface;=20key=5Fprefix=E2=86=92prefix?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Now that the shipped capture clients use keyless Authorization: Bearer (+ PL fingerprint) and :5015 is loopback-only, remove the legacy key-in-URL observer-ingest surface end to end. - Drop the `/X/` route forms on ingest_upload, ingest_manifest, ingest_manifest_day, ingest_event, ingest_segments, and delete_source; the keyless siblings authenticate via PL fingerprint then Bearer. - Delete the dead ingest_transfer route + its require_login exempt entry and structural-gate test entry (zero production callers). - Strip the url_key auth fallback from resolve_observer_identity / _get_auth_key: identity is PL-fingerprint-first then Bearer, with no third fallback. Absent credentials now return auth_required (not the misleading auth_key_invalid). - Collapse the DL arm out of `transfer.py`: remove the URL/key send client, its orphaned helpers, and the DL constants; `transfer send` is PL peer-label only. _normalize_url is retained (export.py imports it). - Rename the observer wire field key_prefix→prefix on /api/list, /api/create, and /init/observers (shared _serialize_observer), the /api/list sort tie-break, and both workspace.html readers. URL route params, internal log fields, and the import-app key_prefix are unchanged. - Document E2 (peer = provenance, not a behavioral role) and E3 (observer-ingest vs import-ingest are deliberately separate) in code comments. The observer callosum SSE route stays keyed (/app/observer//callosum) so shipped Bearer-capable Linux/macOS chat bridges keep working; its segment is now vestigial (unused for auth). The keyless /app/observer/callosum conversion is a deferred coordinated-client follow-up. --- solstone/apps/health/workspace.html | 2 +- solstone/apps/import/ingest.py | 3 + solstone/apps/observer/routes.py | 124 +---- .../apps/observer/tests/test_callosum_sse.py | 45 +- .../observer/tests/test_resolve_identity.py | 12 +- solstone/apps/observer/tests/test_routes.py | 466 +++--------------- solstone/apps/observer/utils.py | 10 +- solstone/apps/observer/workspace.html | 8 +- solstone/convey/root.py | 1 - solstone/observe/transfer.py | 314 +----------- solstone/think/link/auth.py | 13 +- solstone/think/link/join_cli.py | 5 +- tests/baselines/api/observer/ingest-day.json | 6 +- tests/link/test_certless_structural_gate.py | 1 - tests/test_init.py | 6 +- tests/test_transfer.py | 292 +---------- tests/test_transfer_pl.py | 68 +-- tests/verify_api.py | 2 +- 18 files changed, 169 insertions(+), 1209 deletions(-) diff --git a/solstone/apps/health/workspace.html b/solstone/apps/health/workspace.html index 26ae912f6..58819222c 100644 --- a/solstone/apps/health/workspace.html +++ b/solstone/apps/health/workspace.html @@ -3175,7 +3175,7 @@ const nameEl = document.createElement('span'); nameEl.className = 'registered-observer-name'; - nameEl.textContent = observer.name || observer.key_prefix || 'observer'; + nameEl.textContent = observer.name || observer.prefix || 'observer'; row.appendChild(nameEl); const labelEl = document.createElement('span'); diff --git a/solstone/apps/import/ingest.py b/solstone/apps/import/ingest.py index fb205185f..5e1b56a7d 100644 --- a/solstone/apps/import/ingest.py +++ b/solstone/apps/import/ingest.py @@ -103,6 +103,9 @@ from .facet_ingest import process_facet def register_ingest_routes(bp) -> None: + # Import ingest is the bulk producer stream for peer/export segments plus + # entities, facets, imports, and config; it is deliberately separate from + # live observer ingest. @bp.route("/journal//ingest/segments", methods=["POST"]) @require_journal_source def ingest_segments(key_prefix: str): diff --git a/solstone/apps/observer/routes.py b/solstone/apps/observer/routes.py index 00cd6c349..b4c2567f4 100644 --- a/solstone/apps/observer/routes.py +++ b/solstone/apps/observer/routes.py @@ -6,7 +6,6 @@ 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 @@ -186,7 +185,7 @@ def _serialize_observer(observer: dict[str, Any], current_now: int) -> dict[str, ) key_prefix = observer_filename_prefix(observer) return { - "key_prefix": key_prefix, + "prefix": key_prefix, "name": observer.get("name", ""), "created_at": observer.get("created_at", 0), "last_seen": observer.get("last_seen"), @@ -219,7 +218,7 @@ def api_list() -> Any: group_order[observer.get("group", "inactive")], 1 if observer.get("last_seen") is None else 0, -(observer.get("last_seen") or 0), - observer.get("key_prefix", ""), + observer.get("prefix", ""), ) ) @@ -253,10 +252,14 @@ def api_list() -> Any: # The feed does NOT filter events. # The feed does NOT redact fields (v1 trust call; same trust boundary as the existing # Convey SSE bridge — observers are inside it). +# Keyless `/app/observer/callosum` is deferred until solstone-linux and +# solstone-macos ship clients that no longer hardcode this keyed URL. Shipped +# 0.3.0 / 1.3.x clients use `/app/observer//callosum`; the `` segment +# is retained transitionally and is no longer used for auth. @observer_bp.route(_OBSERVER_CALLOSUM_SSE_RULE, methods=["GET"]) def callosum_sse(key: str) -> Any: """Stream Callosum events to an authenticated observer process.""" - observer, key_prefix, error = resolve_observer_identity(key) + observer, key_prefix, error = resolve_observer_identity() if error is not None: return error auth_key = observer.get("key") @@ -322,8 +325,7 @@ def api_create() -> Any: (apps/observer/workspace.html). Auto-registering observer clients use POST /app/observer/register instead, which takes a self-descriptor and locks a stream identity onto the record. This route is kept for the - human-facing management flow; it always mints and returns ``key_prefix`` - (vs /register's ``prefix``). + human-facing management flow; it always mints and returns ``prefix``. """ data = request.get_json(force=True) if request.is_json else {} name = data.get("name", "").strip() @@ -367,7 +369,7 @@ def api_create() -> Any: return jsonify( { "key": key, - "key_prefix": key[:8], + "prefix": key[:8], "name": name, "ingest_url": ingest_url, "protocol_version": protocol.OBSERVER_PROTOCOL_VERSION, @@ -613,10 +615,9 @@ def _save_to_failed( @observer_bp.route("/source/", methods=["DELETE"]) -@observer_bp.route("/source//", methods=["DELETE"]) -def delete_source(stream: str, key: str | None = None) -> Any: +def delete_source(stream: str) -> Any: """Delete an allowed source stream for an authenticated observer.""" - observer, key_prefix, error = resolve_observer_identity(key) + observer, key_prefix, error = resolve_observer_identity() if error is not None: return error @@ -858,10 +859,11 @@ def _process_ingest_files( @observer_bp.route("/ingest", methods=["POST"]) -@observer_bp.route("/ingest/", methods=["POST"]) -def ingest_upload(key: str | None = None) -> Any: +def ingest_upload() -> Any: """Receive file uploads from observer. + Observer ingest is the live, single capture-segment stream from one observer. + Expects multipart form with: - segment: Segment key (HHMMSS_LEN) - day: Day string (YYYYMMDD) @@ -878,7 +880,7 @@ def ingest_upload(key: str | None = None) -> Any: - "duplicate": All files already received (no processing triggered) - "collision": New segment saved with adjusted key (directory conflict) """ - observer, key_prefix, error = resolve_observer_identity(key) + observer, key_prefix, error = resolve_observer_identity() if error is not None: return error @@ -993,86 +995,10 @@ def ingest_upload(key: str | None = None) -> Any: return jsonify(body), status -@observer_bp.route("/ingest//transfer", methods=["POST"]) -def ingest_transfer(key: str) -> Any: - """Receive transferred file uploads from another solstone instance.""" - observer, key_prefix, error = resolve_observer_identity(key) - if error is not None: - return error - - 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 error_response(MISSING_REQUIRED_FIELD, detail="Missing segment") - if not day: - return error_response(MISSING_REQUIRED_FIELD, detail="Missing day") - if not stream: - return error_response(MISSING_REQUIRED_FIELD, detail="Missing stream") - if not re.match(r"^\d{6}_\d+$", segment): - return error_response( - INVALID_SEGMENT_OR_STREAM, - detail="Invalid segment format", - ) - if not re.match(r"^\d{8}$", day): - return error_response(INVALID_DAY, detail="Invalid day format") - if not re.match(r"^[a-z0-9][a-z0-9._-]*$", stream): - return error_response( - INVALID_SEGMENT_OR_STREAM, - detail="Invalid stream format", - ) - - files = request.files.getlist("files") - if not files: - return error_response(INGEST_NO_FILES, detail="No files uploaded") - - 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"]) -@observer_bp.route("/ingest//manifest", methods=["GET"]) -def ingest_manifest(key: str | None = None) -> Any: +def ingest_manifest() -> Any: """List available manifest days for an observer.""" - _observer, key_prefix, error = resolve_observer_identity(key) + _observer, key_prefix, error = resolve_observer_identity() if error is not None: return error @@ -1094,10 +1020,9 @@ def ingest_manifest(key: str | None = None) -> Any: @observer_bp.route("/ingest/manifest/", methods=["GET"]) -@observer_bp.route("/ingest//manifest/", methods=["GET"]) -def ingest_manifest_day(day: str, key: str | None = None) -> Any: +def ingest_manifest_day(day: str) -> Any: """Return a transfer manifest for all segments on a given day.""" - _observer, _key_prefix, error = resolve_observer_identity(key) + _observer, _key_prefix, error = resolve_observer_identity() if error is not None: return error @@ -1130,8 +1055,7 @@ def ingest_manifest_day(day: str, key: str | None = None) -> Any: @observer_bp.route("/ingest/event", methods=["POST"]) -@observer_bp.route("/ingest//event", methods=["POST"]) -def ingest_event(key: str | None = None) -> Any: +def ingest_event() -> Any: """Receive events from observer and relay to local Callosum. Expects JSON body with: @@ -1139,7 +1063,7 @@ def ingest_event(key: str | None = None) -> Any: - event: Event name - ...additional fields """ - observer, _key_prefix, error = resolve_observer_identity(key) + observer, _key_prefix, error = resolve_observer_identity() if error is not None: return error @@ -1199,8 +1123,7 @@ def _respond_observer_segments(items: list[dict], *, client_pv: int) -> Any: @observer_bp.route("/ingest/segments/") -@observer_bp.route("/ingest//segments/") -def ingest_segments(day: str, key: str | None = None) -> Any: +def ingest_segments(day: str) -> Any: """List uploaded segments for a day with file verification. Returns JSON array of segments with file status: @@ -1210,9 +1133,8 @@ def ingest_segments(day: str, key: str | None = None) -> Any: Args: day: Day string (YYYYMMDD) - key: Observer authentication key (from URL path, legacy) """ - observer, key_prefix, error = resolve_observer_identity(key) + observer, key_prefix, error = resolve_observer_identity() if error is not None: return error diff --git a/solstone/apps/observer/tests/test_callosum_sse.py b/solstone/apps/observer/tests/test_callosum_sse.py index 47d805b19..4cb1390c8 100644 --- a/solstone/apps/observer/tests/test_callosum_sse.py +++ b/solstone/apps/observer/tests/test_callosum_sse.py @@ -56,7 +56,7 @@ def _create_observer(env, name: str = "sse-test") -> tuple[str, str]: ) assert resp.status_code == 200 data = resp.get_json() - return data["key"], data["key_prefix"] + return data["key"], data["prefix"] def _route_for(key: str) -> str: @@ -103,11 +103,16 @@ def test_callosum_sse_missing_key_returns_401(observer_env): ) -def test_callosum_sse_unknown_key_returns_401(observer_env): +def test_callosum_sse_path_key_without_bearer_returns_401(observer_env): env = observer_env() - resp = env.client.get(_route_for("unknown-key"), buffered=False) + key, _ = _create_observer(env) + resp = env.client.get(_route_for(key), buffered=False) assert resp.status_code == 401 - _assert_reason(resp, reason_code="auth_key_invalid", detail="Invalid key") + _assert_reason( + resp, + reason_code="auth_required", + detail="Authorization required", + ) def test_callosum_sse_revoked_key_returns_403(observer_env): @@ -116,7 +121,11 @@ def test_callosum_sse_revoked_key_returns_403(observer_env): revoke = env.client.delete(f"/app/observer/api/{key_prefix}") assert revoke.status_code == 200 - resp = env.client.get(_route_for(key), buffered=False) + resp = env.client.get( + _route_for("vestigial-path-key"), + headers={"Authorization": f"Bearer {key}"}, + buffered=False, + ) assert resp.status_code == 403 _assert_reason( resp, @@ -133,7 +142,11 @@ def test_callosum_sse_disabled_key_returns_403(observer_env): observer["enabled"] = False assert save_observer(observer) - resp = env.client.get(_route_for(key), buffered=False) + resp = env.client.get( + _route_for("vestigial-path-key"), + headers={"Authorization": f"Bearer {key}"}, + buffered=False, + ) assert resp.status_code == 403 _assert_reason( resp, @@ -142,7 +155,7 @@ def test_callosum_sse_disabled_key_returns_403(observer_env): ) -def test_callosum_sse_bearer_header_overrides_path_key(observer_env): +def test_callosum_sse_bearer_header_authenticates(observer_env): env = observer_env() valid_key, _ = _create_observer(env, "valid-sse") bogus_key = "bogus-key" @@ -171,7 +184,11 @@ def test_callosum_sse_success_content_type(observer_env): env = observer_env() key, _ = _create_observer(env) - resp = env.client.get(_route_for(key), buffered=False) + resp = env.client.get( + _route_for("vestigial-path-key"), + headers={"Authorization": f"Bearer {key}"}, + buffered=False, + ) try: assert resp.status_code == 200 assert resp.content_type.startswith("text/event-stream") @@ -182,7 +199,11 @@ def test_callosum_sse_success_content_type(observer_env): def test_callosum_sse_round_trip_payload(observer_env): env = observer_env() key, key_prefix = _create_observer(env) - resp = env.client.get(_route_for(key), buffered=False) + resp = env.client.get( + _route_for("vestigial-path-key"), + headers={"Authorization": f"Bearer {key}"}, + buffered=False, + ) try: assert resp.status_code == 200 assert convey_bridge.subscription_count(key_prefix) == 1 @@ -268,7 +289,11 @@ def test_callosum_sse_heartbeat(observer_env, monkeypatch): key, _ = _create_observer(env) monkeypatch.setattr(routes_module, "_SSE_HEARTBEAT_SECONDS", 0.01) - resp = env.client.get(_route_for(key), buffered=False) + resp = env.client.get( + _route_for("vestigial-path-key"), + headers={"Authorization": f"Bearer {key}"}, + buffered=False, + ) try: assert resp.status_code == 200 assert _next_chunk(resp) == ": heartbeat\n\n" diff --git a/solstone/apps/observer/tests/test_resolve_identity.py b/solstone/apps/observer/tests/test_resolve_identity.py index 6299cd89b..0285397dd 100644 --- a/solstone/apps/observer/tests/test_resolve_identity.py +++ b/solstone/apps/observer/tests/test_resolve_identity.py @@ -57,16 +57,14 @@ def test_resolve_dl_success_from_bearer(app_env): assert prefix == DL_KEY[:8] -def test_resolve_dl_bearer_key_wins_over_url_key(app_env): +def test_resolve_dl_uses_bearer_key(app_env): header_key = "headerkey123456789" - url_key = "urlkey123456789" save_observer({"key": header_key, "name": "header", "enabled": True, "stats": {}}) - save_observer({"key": url_key, "name": "url", "enabled": True, "stats": {}}) with app_env.test_request_context( headers={"Authorization": f"Bearer {header_key}"} ): - observer, prefix, error = resolve_observer_identity(url_key=url_key) + observer, prefix, error = resolve_observer_identity() assert error is None assert observer["name"] == "header" @@ -125,7 +123,7 @@ def test_resolve_pl_success(app_env): with app_env.test_request_context(): g.identity = _pl_identity(FINGERPRINT) - observer, prefix, error = resolve_observer_identity("ignored-route-key") + observer, prefix, error = resolve_observer_identity() assert error is None assert observer["name"] == "observer" @@ -182,12 +180,12 @@ def test_resolve_pl_disabled(app_env): assert payload["reason_code"] == "feature_unavailable" -def test_resolve_pl_does_not_require_url_key(app_env): +def test_resolve_pl_does_not_require_bearer(app_env): mint_pl_observer_record(FINGERPRINT, "observer", "2026-04-20T00:00:00Z") with app_env.test_request_context("/app/observer/ingest/not-the-fingerprint/event"): g.identity = _pl_identity(FINGERPRINT) - observer, prefix, error = resolve_observer_identity("not-the-fingerprint") + observer, prefix, error = resolve_observer_identity() assert error is None assert prefix == "c" * 16 diff --git a/solstone/apps/observer/tests/test_routes.py b/solstone/apps/observer/tests/test_routes.py index c76bff576..272a2eccc 100644 --- a/solstone/apps/observer/tests/test_routes.py +++ b/solstone/apps/observer/tests/test_routes.py @@ -9,8 +9,6 @@ import io import json from unittest.mock import MagicMock -import requests - import solstone.apps.observer.routes as routes_module import solstone.convey.bridge as convey_bridge from solstone.apps.observer.routes import ( @@ -259,7 +257,8 @@ def test_api_create_observer(observer_env): assert "key" in data assert len(data["key"]) > 32 # 256 bits = 43 base64 chars - assert data["key_prefix"] == data["key"][:8] + assert data["prefix"] == data["key"][:8] + assert "key_prefix" not in data assert data["name"] == "test-laptop" assert data["ingest_url"] == "/app/observer/ingest" assert data["key"] not in data["ingest_url"] @@ -299,7 +298,7 @@ def test_api_list_shows_created_observer(observer_env): content_type="application/json", ) assert resp.status_code == 200 - key_prefix = resp.get_json()["key_prefix"] + key_prefix = resp.get_json()["prefix"] # List should show it payload = _api_list_payload(env) @@ -307,7 +306,8 @@ def test_api_list_shows_created_observer(observer_env): assert len(observers) == 1 assert payload["thresholds"] == {"active_ms": 30000, "stale_ms": 120000} - assert observers[0]["key_prefix"] == key_prefix + assert observers[0]["prefix"] == key_prefix + assert "key_prefix" not in observers[0] assert observers[0]["name"] == "my-observer" assert observers[0]["enabled"] is True assert observers[0]["stats"]["segments_received"] == 0 @@ -327,7 +327,7 @@ def test_api_list_includes_last_chat_request_at(observer_env): content_type="application/json", ) assert resp.status_code == 200 - key_prefix = resp.get_json()["key_prefix"] + key_prefix = resp.get_json()["prefix"] handle = convey_bridge.register_sse_subscriber(key_prefix) try: convey_bridge._broadcast_to_sse_clients( @@ -339,7 +339,7 @@ def test_api_list_includes_last_chat_request_at(observer_env): with convey_bridge._SSE_LOCK: convey_bridge._SSE_LAST_CHAT_REQUEST_AT_BY_KEY.pop(key_prefix, None) - assert observers[0]["key_prefix"] == key_prefix + assert observers[0]["prefix"] == key_prefix assert observers[0]["last_chat_request_at"] == 9876 @@ -353,7 +353,7 @@ def test_api_delete_observer(observer_env): json={"name": "to-revoke"}, content_type="application/json", ) - key_prefix = resp.get_json()["key_prefix"] + key_prefix = resp.get_json()["prefix"] # Revoke it resp = env.client.delete(f"/app/observer/api/{key_prefix}") @@ -363,7 +363,7 @@ def test_api_delete_observer(observer_env): # List should still show it, but marked as revoked observers = _api_list_observers(env) assert len(observers) == 1 - assert observers[0]["key_prefix"] == key_prefix + assert observers[0]["prefix"] == key_prefix assert observers[0]["revoked"] is True assert observers[0]["revoked_at"] is not None assert observers[0]["state"] == "revoked" @@ -409,7 +409,7 @@ def test_api_delete_dl_observer_does_not_touch_authorized_clients(observer_env): json={"name": "dl-delete"}, content_type="application/json", ) - key_prefix = resp.get_json()["key_prefix"] + key_prefix = resp.get_json()["prefix"] fingerprint = "sha256:" + ("e" * 64) AuthorizedClients(authorized_clients_path()).add( fingerprint, @@ -495,8 +495,8 @@ def test_api_list_sorts_by_group_and_last_seen(observer_env, monkeypatch): ] -def test_api_list_tie_breaks_by_key_prefix(observer_env, monkeypatch): - """Observers with the same last_seen sort by key_prefix ascending.""" +def test_api_list_tie_breaks_by_prefix(observer_env, monkeypatch): + """Observers with the same last_seen sort by prefix ascending.""" env = observer_env() fixed_now = 3_000_000 monkeypatch.setattr(routes_module, "now_ms", lambda: fixed_now) @@ -515,10 +515,11 @@ def test_api_list_tie_breaks_by_key_prefix(observer_env, monkeypatch): ) observers = _api_list_observers(env) - assert [observer["key_prefix"] for observer in observers] == [ + assert [observer["prefix"] for observer in observers] == [ "aaaa0000", "bbbb0000", ] + assert all("key_prefix" not in observer for observer in observers) assert all(observer["state"] == "connected" for observer in observers) assert all(observer["group"] == "active" for observer in observers) assert all( @@ -636,13 +637,14 @@ def test_delete_source_hard_pin_rejects_other_stream(observer_env): content_type="application/json", ) key = create_resp.get_json()["key"] + headers = {"Authorization": f"Bearer {key}"} other_seg = _plant_source_segment( env, stream="import.audio", segment="130000_300", ) - resp = env.client.delete(f"/app/observer/source/import.audio/{key}") + resp = env.client.delete("/app/observer/source/import.audio", headers=headers) assert resp.status_code == 400 assert resp.get_json()["detail"] == "Only known source streams can be deleted" assert other_seg.exists() @@ -653,7 +655,8 @@ def test_delete_source_hard_pin_rejects_other_stream(observer_env): segment="140000_300", ) resp = env.client.delete( - f"/app/observer/source/import.share/{key}", + "/app/observer/source/import.share", + headers=headers, data={"stream": "import.audio"}, ) assert resp.status_code == 400 @@ -661,7 +664,8 @@ def test_delete_source_hard_pin_rejects_other_stream(observer_env): assert share_seg.exists() resp = env.client.delete( - f"/app/observer/source/import.share/{key}", + "/app/observer/source/import.share", + headers=headers, data={"meta": json.dumps({"stream": "import.audio"})}, ) assert resp.status_code == 400 @@ -680,7 +684,10 @@ def test_delete_source_happy_path(observer_env): key = create_resp.get_json()["key"] seg_dir = _plant_source_segment(env) - resp = env.client.delete(f"/app/observer/source/import.share/{key}") + resp = env.client.delete( + "/app/observer/source/import.share", + headers={"Authorization": f"Bearer {key}"}, + ) assert resp.status_code == 200 receipt = resp.get_json() @@ -718,7 +725,10 @@ def test_delete_source_location_happy_path(observer_env): key = create_resp.get_json()["key"] seg_dir = _plant_location_segment(env) - resp = env.client.delete(f"/app/observer/source/location/{key}") + resp = env.client.delete( + "/app/observer/source/location", + headers={"Authorization": f"Bearer {key}"}, + ) assert resp.status_code == 200 receipt = resp.get_json() @@ -737,6 +747,7 @@ def test_delete_source_path_wins_over_candidate(observer_env): content_type="application/json", ) key = create_resp.get_json()["key"] + headers = {"Authorization": f"Bearer {key}"} form_location_seg = _plant_location_segment( env, day="20250105", @@ -749,7 +760,8 @@ def test_delete_source_path_wins_over_candidate(observer_env): ) form_resp = env.client.delete( - f"/app/observer/source/location/{key}", + "/app/observer/source/location", + headers=headers, data={"stream": "import.share"}, ) @@ -769,7 +781,8 @@ def test_delete_source_path_wins_over_candidate(observer_env): ) meta_resp = env.client.delete( - f"/app/observer/source/location/{key}", + "/app/observer/source/location", + headers=headers, data={"meta": json.dumps({"stream": "import.share"})}, ) @@ -1011,44 +1024,6 @@ def test_ingest_event_relay(observer_env): assert resp.get_json()["status"] == "ok" -def test_ingest_event_pl_ignores_url_key(observer_env, monkeypatch): - env = observer_env() - other_key = _save_test_observer( - "deadbeef", - "other-dl", - created_at=100, - last_seen=None, - ) - mint_pl_observer_record( - fingerprint=PL_FINGERPRINT, - device_label="pl-event", - paired_at="2026-05-20T00:00:00Z", - ) - emitted: list[tuple[str, str, dict]] = [] - monkeypatch.setattr( - routes_module, - "emit", - lambda tract, event, **kwargs: emitted.append((tract, event, kwargs)), - ) - - resp = env.client.post( - f"/app/observer/ingest/{other_key[:8]}/event", - environ_overrides={"pl.identity": _pl_identity()}, - json={"tract": "observe", "event": "status", "mode": "screencast"}, - content_type="application/json", - ) - - assert resp.status_code == 200 - assert resp.get_json()["status"] == "ok" - assert emitted == [ - ( - "observe", - "status", - {"mode": "screencast", "observer": "pl-event"}, - ) - ] - - def test_dl_and_pl_observers_coexist_and_ingest(observer_env): env = observer_env() dl_resp = env.client.post( @@ -1164,7 +1139,7 @@ def test_ingest_revoked_key(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] resp = env.client.delete(f"/app/observer/api/{key_prefix}") assert resp.status_code == 200 @@ -1195,7 +1170,7 @@ def test_keyless_ingest_bearer_rejects_revoked_and_disabled_keys(observer_env): revoked_data = resp.get_json() revoked_key = revoked_data["key"] - resp = env.client.delete(f"/app/observer/api/{revoked_data['key_prefix']}") + resp = env.client.delete(f"/app/observer/api/{revoked_data['prefix']}") assert resp.status_code == 200 resp = env.client.post( @@ -1263,7 +1238,7 @@ def test_ingest_event_revoked_key(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] resp = env.client.delete(f"/app/observer/api/{key_prefix}") assert resp.status_code == 200 @@ -1291,7 +1266,7 @@ def test_api_get_key(observer_env): ) create_data = resp.get_json() key = create_data["key"] - key_prefix = create_data["key_prefix"] + key_prefix = create_data["prefix"] # Get the key resp = env.client.get(f"/app/observer/api/{key_prefix}/key") @@ -1321,7 +1296,7 @@ def test_mint_responses_protocol_version_single_source_and_keyless_unconditional assert create_data["protocol_version"] == 99 assert create_data["ingest_url"] == "/app/observer/ingest" - resp = env.client.get(f"/app/observer/api/{create_data['key_prefix']}/key") + resp = env.client.get(f"/app/observer/api/{create_data['prefix']}/key") assert resp.status_code == 200 key_data = resp.get_json() assert key_data["protocol_version"] == 99 @@ -1347,7 +1322,7 @@ def test_api_get_key_revoked(observer_env): content_type="application/json", ) create_data = resp.get_json() - key_prefix = create_data["key_prefix"] + key_prefix = create_data["prefix"] env.client.delete(f"/app/observer/api/{key_prefix}") @@ -1369,7 +1344,7 @@ def test_api_get_key_audit_log(observer_env): content_type="application/json", ) create_data = resp.get_json() - key_prefix = create_data["key_prefix"] + key_prefix = create_data["prefix"] with patch("solstone.apps.observer.routes.log_app_action") as mock_log: resp = env.client.get(f"/app/observer/api/{key_prefix}/key") @@ -1607,7 +1582,7 @@ def test_ingest_creates_sync_history(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] # Upload a file test_data = b"test audio content for history" @@ -1663,7 +1638,7 @@ def test_ingest_history_with_collision(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] # Create conflicting segment directory under the stream day_dir = _day_dir(env) @@ -1812,41 +1787,6 @@ def test_segments_endpoint_lists_uploads(observer_env): ) # Original name preserved -def test_legacy_url_key_ingest_and_segments_still_sync(observer_env): - env = observer_env() - - resp = env.client.post( - "/app/observer/api/create", - json={"name": "legacy-url-key-test"}, - content_type="application/json", - ) - key = resp.get_json()["key"] - - test_data = b"legacy url key content" - resp = env.client.post( - f"/app/observer/ingest/{key}", - data={ - "day": "20250103", - "segment": "120000_300", - "files": (io.BytesIO(test_data), "120000_300_audio.flac"), - }, - ) - assert resp.status_code == 200 - - resp = env.client.get(f"/app/observer/ingest/{key}/segments/20250103") - assert resp.status_code == 200 - data = resp.get_json() - - assert isinstance(data, list) - assert len(data) == 1 - segment = data[0] - assert segment["key"] == "120000_300" - assert len(segment["files"]) == 1 - - file_info = segment["files"][0] - assert file_info["status"] == "present" - - def test_segments_endpoint_v2_empty(observer_env): """Test v2 segments endpoint returns collection envelope for no uploads.""" env = observer_env() @@ -1949,20 +1889,6 @@ def test_protocol_version_single_source(observer_env, monkeypatch): assert resp.status_code == 200 assert isinstance(resp.get_json(), list) - from solstone.observe import transfer - - session = MagicMock(spec=requests.Session) - response = MagicMock() - response.status_code = 200 - response.json.return_value = [] - session.get.return_value = response - - transfer._query_remote_segments(session, "http://x", "20250103") - - assert ( - session.get.call_args.kwargs["headers"]["X-Solstone-Protocol-Version"] == "99" - ) - def test_segments_endpoint_shows_collision(observer_env): """Test segments endpoint shows collision info.""" @@ -2145,7 +2071,7 @@ def test_segments_endpoint_revoked_key(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] env.client.delete(f"/app/observer/api/{key_prefix}") @@ -2228,7 +2154,7 @@ def test_segments_endpoint_shows_observed_status(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] # Upload a file test_data = b"test audio content" @@ -2280,7 +2206,7 @@ def test_api_list_includes_segments_observed_stat(observer_env): content_type="application/json", ) data = resp.get_json() - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] # Initially no segments_observed data = _api_list_observers(env) @@ -2523,7 +2449,7 @@ def test_ingest_partial_match_logged_in_history(observer_env): ) data = resp.get_json() key = data["key"] - key_prefix = data["key_prefix"] + key_prefix = data["prefix"] audio_data = b"test audio for partial log" @@ -2716,286 +2642,6 @@ def test_ingest_stream_qualifier_preserved(observer_env): assert not (_day_dir(env) / "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 = _day_dir(env) / "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()["detail"] == "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()["detail"] == "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 = _day_dir(env) / "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 solstone.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 solstone.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 solstone.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() @@ -3018,7 +2664,10 @@ def test_manifest_day_listing(observer_env): ) assert resp.status_code == 200 - resp = env.client.get(f"/app/observer/ingest/{key}/manifest") + resp = env.client.get( + "/app/observer/ingest/manifest", + headers={"Authorization": f"Bearer {key}"}, + ) assert resp.status_code == 200 assert resp.get_json() == {"days": {"20250103": {"segments": 1}}} @@ -3035,11 +2684,12 @@ def test_manifest_per_day(observer_env): key = resp.get_json()["key"] resp = env.client.post( - f"/app/observer/ingest/{key}/transfer", + "/app/observer/ingest", + headers={"Authorization": f"Bearer {key}"}, data={ "day": "20250103", "segment": "120000_300", - "stream": "remote.host", + "meta": json.dumps({"stream": "remote.host"}), "files": [ (io.BytesIO(b"audio bytes"), "audio.flac"), (io.BytesIO(b"screen bytes"), "screen.webm"), @@ -3048,7 +2698,10 @@ def test_manifest_per_day(observer_env): ) assert resp.status_code == 200 - resp = env.client.get(f"/app/observer/ingest/{key}/manifest/20250103") + resp = env.client.get( + "/app/observer/ingest/manifest/20250103", + headers={"Authorization": f"Bearer {key}"}, + ) assert resp.status_code == 200 data = resp.get_json() @@ -3059,7 +2712,8 @@ def test_manifest_per_day(observer_env): assert "remote.host/120000_300" in data["segments"] files = data["segments"]["remote.host/120000_300"]["files"] - assert len(files) == 2 + names = {file_info["name"] for file_info in files} + assert {"audio.flac", "screen.webm"}.issubset(names) for file_info in files: assert set(file_info) == {"name", "sha256", "size"} assert len(file_info["sha256"]) == 64 @@ -3069,5 +2723,5 @@ 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") + resp = env.client.get("/app/observer/ingest/manifest") assert resp.status_code == 401 diff --git a/solstone/apps/observer/utils.py b/solstone/apps/observer/utils.py index c5030d26d..30a33fae2 100644 --- a/solstone/apps/observer/utils.py +++ b/solstone/apps/observer/utils.py @@ -361,7 +361,7 @@ def find_observer_by_name(name: str) -> dict | None: return ObserverRegistry.singleton().by_name(name) -def _get_auth_key(url_key: str | None = None) -> str | None: +def _get_auth_key() -> str | None: from flask import request auth = request.headers.get("Authorization", "") @@ -369,7 +369,7 @@ def _get_auth_key(url_key: str | None = None) -> str | None: bearer = auth[7:].strip() if bearer: return bearer - return url_key or None + return None def _identity_fingerprint() -> str | None: @@ -394,8 +394,8 @@ def _check_observer_enabled(observer: dict): return None -# Observer still resolves through ObserverRegistry; minted pairings use header-Bearer auth, with legacy key-in-URL fallback deferring require_ingest_identity until URL-key auth is retired. -def resolve_observer_identity(url_key: str | None = None): +def resolve_observer_identity(): + """Resolve an observer from PL identity first, then Authorization Bearer.""" fingerprint = _identity_fingerprint() if fingerprint is not None: observer = load_observer_by_fingerprint(fingerprint) @@ -406,7 +406,7 @@ def resolve_observer_identity(url_key: str | None = None): return None, None, error return observer, observer["filename_prefix"], None - auth_key = _get_auth_key(url_key) + auth_key = _get_auth_key() if not auth_key: return None, None, _auth_failure() observer = load_observer(auth_key) diff --git a/solstone/apps/observer/workspace.html b/solstone/apps/observer/workspace.html index 85a6a542a..432e679da 100644 --- a/solstone/apps/observer/workspace.html +++ b/solstone/apps/observer/workspace.html @@ -707,7 +707,7 @@ function observerCardHTML(observer, labels = {}) { const ariaLabel = [observer.name, label, liveLabel].filter(Boolean).join(', '); return ` -
+
${escapeHtml(observer.name)} ${safeLabel} @@ -715,9 +715,9 @@ function observerCardHTML(observer, labels = {}) {
${statsHTML(observer, statusClass)}
- ${statusClass === 'stale' && !observer.revoked ? `` : ''} - ${observer.revoked ? '' : ``} - ${observer.revoked ? '' : ``} + ${statusClass === 'stale' && !observer.revoked ? `` : ''} + ${observer.revoked ? '' : ``} + ${observer.revoked ? '' : ``}
`; diff --git a/solstone/convey/root.py b/solstone/convey/root.py index e594fbbd2..334ebafa1 100644 --- a/solstone/convey/root.py +++ b/solstone/convey/root.py @@ -110,7 +110,6 @@ def require_login() -> Any: "app:observer.ingest_upload", "app:observer.ingest_event", "app:observer.ingest_segments", - "app:observer.ingest_transfer", "app:observer.ingest_manifest", "app:observer.ingest_manifest_day", "app:observer.register", diff --git a/solstone/observe/transfer.py b/solstone/observe/transfer.py index 5c301b997..e4ea517b9 100644 --- a/solstone/observe/transfer.py +++ b/solstone/observe/transfer.py @@ -9,12 +9,9 @@ observation segments between solstone instances. Usage: journal transfer export --day YYYYMMDD [--output PATH] journal transfer import --archive PATH [--dry-run] - journal transfer send --to URL --key KEY [--day YYYYMMDD] [--dry-run] journal transfer send --to LABEL [--day YYYYMMDD] [--dry-run] -On the RECEIVING host (the machine you are sending TO), run -`journal observer create ` to generate an observer API key, then pass it -as `--key`. +Use `journal transfer send --to LABEL` to send segments to a paired peer. """ from __future__ import annotations @@ -33,9 +30,6 @@ from datetime import datetime, timedelta from pathlib import Path from typing import Any -import requests - -from solstone.observe import protocol from solstone.observe.peer_lookup import PeerInfo, PeerLookupError, resolve_peer from solstone.observe.pl_http import PlHttpSession from solstone.think.callosum import callosum_send @@ -56,16 +50,6 @@ from solstone.think.utils import ( from .utils import compute_file_sha256, find_available_segment -OBSERVER_KEY_HINT = ( - "On the RECEIVING host (the machine you are sending TO), run " - "`journal observer create ` to generate an observer API key, then " - "pass it as `--key`." -) - -AUTH_INVALID_OBSERVER_KEY = ( - "Authentication failed: invalid or missing observer API key. " + OBSERVER_KEY_HINT -) - logger = logging.getLogger(__name__) # Archive manifest version @@ -448,26 +432,12 @@ def _normalize_url(to: str) -> str: return f"https://{to}" -def _is_url_destination(to: str) -> bool: - return to.startswith(("http://", "https://")) - - def _resolve_destination( parser: argparse.ArgumentParser, to: str, - key: str | None, -) -> tuple[str, str | None, PeerInfo | None]: - if _is_url_destination(to): - if key is None: - parser.error("'--to ' requires '--key '") - return ("dl", _normalize_url(to), None) - if key is not None: - parser.error( - "'--key' is only valid with '--to '; " - "use '--to