diff --git a/convey/__init__.py b/convey/__init__.py index fae75fe86..ff52a0f79 100644 --- a/convey/__init__.py +++ b/convey/__init__.py @@ -109,8 +109,10 @@ def _migrate_setup_completed() -> None: def create_app(journal: str = "") -> Flask: """Create and configure the Convey Flask application.""" + from think.push.runtime import start_push_runtime from think.voice.runtime import start_voice_runtime + from .push import push_bp from .voice import voice_bp app = Flask( @@ -150,6 +152,9 @@ def create_app(journal: str = "") -> Flask: # Register voice API blueprint app.register_blueprint(voice_bp) + # Register push API blueprint + app.register_blueprint(push_bp) + # Initialize and register app system registry = AppRegistry() registry.discover() @@ -161,6 +166,7 @@ def create_app(journal: str = "") -> Flask: sock = Sock(app) register_websocket(sock) start_voice_runtime(app) + start_push_runtime(app) if journal: state.journal_root = journal diff --git a/convey/push.py b/convey/push.py new file mode 100644 index 000000000..ec784b18b --- /dev/null +++ b/convey/push.py @@ -0,0 +1,127 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Root-level push API.""" + +from __future__ import annotations + +import uuid +from typing import Any + +from flask import Blueprint, jsonify, request +from werkzeug.exceptions import BadRequest + +from think.push import triggers +from think.push.config import is_configured +from think.push.devices import ( + load_devices, + register_device, + remove_device, + status_view, +) +from think.push.dispatch import CATEGORIES, CATEGORY_AGENT_ALERT + +push_bp = Blueprint("push", __name__, url_prefix="/api/push") + + +def _error(message: str, status: int): + return jsonify({"error": message}), status + + +def _optional_json_object() -> tuple[dict[str, Any], Any | None]: + if not request.get_data(cache=True): + return {}, None + try: + data = request.get_json(silent=False) + except BadRequest: + return {}, _error("request body must be valid JSON", 400) + if not isinstance(data, dict): + return {}, _error("request body must be a JSON object", 400) + return data, None + + +def _required_json_object() -> tuple[dict[str, Any], Any | None]: + try: + data = request.get_json(silent=False) + except BadRequest: + return {}, _error("request body must be valid JSON", 400) + if not isinstance(data, dict): + return {}, _error("request body must be a JSON object", 400) + return data, None + + +@push_bp.post("/register") +def register_push_device(): + body, error = _required_json_object() + if error is not None: + return error + token = str(body.get("device_token") or "").strip() + bundle_id = str(body.get("bundle_id") or "").strip() + environment = str(body.get("environment") or "").strip() + platform = str(body.get("platform") or "").strip() + if not token: + return _error("device_token is required", 400) + if not bundle_id: + return _error("bundle_id is required", 400) + if environment not in {"development", "production"}: + return _error("environment must be development or production", 400) + if platform != "ios": + return _error("platform must be ios", 400) + count = register_device( + token="".join(token.split()).lower(), + bundle_id=bundle_id, + environment=environment, + platform=platform, + ) + return jsonify({"registered": True, "device_count": count}) + + +@push_bp.delete("/register") +def unregister_push_device(): + body, error = _required_json_object() + if error is not None: + return error + token = str(body.get("device_token") or "").strip() + if not token: + return _error("device_token is required", 400) + removed = remove_device("".join(token.split()).lower()) + return jsonify({"removed": removed, "device_count": len(load_devices())}) + + +@push_bp.get("/status") +def push_status(): + devices = sorted( + load_devices(), + key=lambda device: int(device.get("registered_at", 0)), + reverse=True, + ) + return jsonify( + { + "configured": is_configured(), + "device_count": len(devices), + "devices": [status_view(device) for device in devices], + } + ) + + +@push_bp.post("/test") +def send_push_test(): + body, error = _optional_json_object() + if error is not None: + return error + if not is_configured(): + return _error("push not configured", 503) + category = body.get("category", CATEGORY_AGENT_ALERT) + if category not in CATEGORIES: + return _error("category must be a known push category", 400) + title = str(body.get("title") or "Push test") + message = str(body.get("body") or "This is a test notification.") + sent, failed = triggers.send_agent_alert( + title=title, + body=message, + context_id=f"push-test-{uuid.uuid4().hex[:12]}", + ) + return jsonify({"sent": sent, "failed": failed}) + + +__all__ = ["push_bp"] diff --git a/docs/design/push.md b/docs/design/push.md new file mode 100644 index 000000000..3c28bce51 --- /dev/null +++ b/docs/design/push.md @@ -0,0 +1,686 @@ +# Wave 3 push server + +## 1. Summary + +Wave 3 ships a root-level push API on the existing Convey server, backed by a dedicated push runtime and APNs dispatch layer. The server surface is a new `convey/push.py` blueprint mounted at `/api/push/*`, following the same root-blueprint pattern as the Wave 2 voice server rather than adding an `apps/push/` package (`convey/voice.py:26-184`, `convey/__init__.py:150-166`, `docs/design/voice-server.md:1-37`, `0a693381 voice: ship Wave 2 voice server (root /api/voice/*, 9-tool sideband)`). The runtime starts from `convey.create_app()`, owns both a `CallosumConnection` listener and an asyncio loop, dispatches Daily Briefing when `cortex.finish` reports `name=="morning_briefing"`, and runs a 60-second periodic check for Pre-Meeting Prep across enabled facets only (`think/callosum.py:245-346`, `think/cortex.py:433-441`, `think/activities.py:877-890`, `think/facets.py:255-261`). The shipped notification categories are `SOLSTONE_DAILY_BRIEFING`, `SOLSTONE_PRE_MEETING_PREP`, and `SOLSTONE_AGENT_ALERT`; the server also defines `SOLSTONE_COMMITMENT_NUDGE` for client forward-compatibility, but no Wave 3 trigger emits it. + +Wave 3 explicitly defers commitment nudges to Wave 3.1 because the current ledger surface does not provide a machine-readable due date. `LedgerItem` has `when: str | None`, no `due` field, and `age_days` is derived from `opened_at`, so any age-threshold proxy would blur “old” and “overdue” in misleading ways (`think/surfaces/types.py:16-31`, `think/surfaces/ledger.py:395-410`). Wave 3 also defers live APNs validation pending Apple Developer enrollment, and it does not cover the native iOS client implementation, which ships in the companion lode. + +Wave 3 is accepted when the mocked round-trip tests pass, the new push paths remain inside their declared write ownership, and the live-validation handoff is explicit enough that enrollment unblocks the final production exercise without design work. The gate stays the same as other repo work: keep layer hygiene clean, keep behavior testable from the Flask surface downward, and make runtime startup and shutdown deterministic (`scripts/check_layer_hygiene.py:38-72`, `scripts/check_layer_hygiene.py:183-240`, `tests/test_voice_runtime.py:19-103`, `tests/test_voice_integration.py:102-149`). + +## 2. Module layout + +Wave 3 adds a small `think/push/` package plus one new root blueprint. The layout deliberately mirrors the voice-server split where that split is sound, and deliberately does not reuse the voice runtime name because voice shutdown still hardcodes `brain.clear_brain_state()` and voice-specific app attachment (`think/voice/runtime.py:44-50`, `think/voice/runtime.py:88-105`). + +| Path | Role | +|---|---| +| `convey/push.py` | Root Flask blueprint at `/api/push/*`. Defines local request validators mirroring `convey/voice.py` (`_error`, `_required_json_object`, `_optional_json_object`) and exposes `POST /register`, `DELETE /register`, `GET /status`, and `POST /test` (`convey/voice.py:26-55`, `convey/voice.py:58-184`). | +| `convey/__init__.py` | Registers `push_bp` beside the other root blueprints and calls `start_push_runtime(app)` during `create_app()`, mirroring the existing voice wire-in. The implementation must keep using `get_journal()` at call time because `state.journal_root` is assigned only after runtime startup (`convey/__init__.py:138-166`). | +| `think/push/__init__.py` | Package marker plus narrow re-export surface for runtime helpers and the public trigger entry point. This keeps callers out of module-private helpers and matches the small public surfaces used elsewhere in `think/voice/` (`think/voice/runtime.py:121-127`). | +| `think/push/config.py` | Journal-scoped config readers for `push.apns_key_path`, `push.apns_key_id`, `push.apns_team_id`, `push.bundle_id`, and `push.environment`. It mirrors the small-reader style of `think/voice/config.py`, but intentionally does not add env-var fallback because push credentials are journal-scoped operational config, not process-scoped ambient state (`think/voice/config.py:17-42`, `think/journal_default.json:35-39`). | +| `think/push/devices.py` | Device store owner. Exposes `load_devices()`, `register_device(token, bundle_id, environment, platform)`, and `remove_device(token)`. It is the sole writer for `journal/config/push_devices.json`, normalizes token input, recovers from malformed stores by treating them as empty with a warning, and rewrites the store atomically on each mutation. | +| `think/push/dispatch.py` | APNs transport owner. Defines the four category constants, the APNs HTTP/2 client using `httpx`, the ES256 JWT signer using `PyJWT`, the 55-minute bearer-token cache, payload builders, and `send(device, payload)` / `send_many(devices, payload)` helpers. `httpx` is already in the repo dependency set, so only `PyJWT` is added (`pyproject.toml:53`). | +| `think/push/triggers.py` | Trigger owner. Defines `handle_briefing_finish(message)`, `check_pre_meeting_prep(now)`, and `send_agent_alert(title, body, context_id)`. This module is the sole writer for `journal/push/nudge_log.jsonl`; there is no separate `log.py`, so dedupe stays adjacent to trigger decisions. Trigger logic reads only enabled facets, reuses `_load_briefing_md(today)` for Daily Briefing, and reuses `load_activity_records(facet, day)` plus `record["source"] == "anticipated"` filtering for Pre-Meeting Prep (`apps/home/routes.py:149-198`, `apps/home/routes.py:305-337`, `think/activities.py:877-890`, `think/activities.py:937-939`, `think/facets.py:255-261`). | +| `think/push/runtime.py` | Dedicated runtime singleton. Exposes `start_push_runtime(app)`, `stop_push_runtime(app)`, and `stop_all_push_runtime()`. It mirrors the voice runtime’s daemon-thread + asyncio-loop + `atexit` pattern, but keeps push lifecycle independent and owns both the callosum listener and the 60-second periodic task (`think/voice/runtime.py:21-109`, `think/callosum.py:254-346`). | + +Deliberate non-changes: + +- No `sol push` top-level CLI. Wave 3 is a root API plus in-process runtime, matching the Wave 2 voice-server precedent rather than adding a separate command surface (`docs/design/voice-server.md:7-37`, `0a693381 voice: ship Wave 2 voice server (root /api/voice/*, 9-tool sideband)`). +- No `schedules.json` entry. The scheduler only understands `hourly`, `daily`, and `weekly`, which is too coarse for a 15-minute pre-meeting reminder (`think/scheduler.py:29-30`, `think/scheduler.py:375-438`). +- No supervisor hook. `think/supervisor.py::supervise()` is the one-second orchestration loop, but Wave 3 keeps push-domain logic out of supervisor and self-contains it in the push runtime (`think/supervisor.py:1311-1371`). + +`push_devices.json` uses this storage model: + +- Top-level JSON object: `{"devices": [...]}`. +- Each device row stores `token`, `bundle_id`, `environment`, `platform`, and `registered_at`. +- Token identity is unique per row. Re-registering the same token updates the row in place and refreshes `registered_at`. +- Dispatch reads all rows, then filters to rows whose `bundle_id`, `environment`, and `platform` match the current push configuration before sending. + +`nudge_log.jsonl` uses this append-only model: + +- One JSON object per successful trigger fire. +- Common fields: `ts`, `category`, `dedupe_key`, `sent`, `failed`. +- Category-specific context: `day` for Daily Briefing, `activity_id` and `facet` for Pre-Meeting Prep, `context_id` for Agent Alert. +- A line is appended only when at least one device send succeeds. Zero-success attempts stay retryable inside the same trigger window. + +## 3. Flow diagrams + +### 3.1 Device register + +```text +iOS client + -> POST /api/push/register + -> convey/push.py validates JSON body + -> think.push.devices.register_device(...) + -> journal/config/push_devices.json rewrite + -> 200 {"registered": true, "device_count": N} +``` + +This mirrors the voice blueprint’s “validate locally, then hand off to the feature module” pattern, but the handoff is synchronous because device registration is just a journal write and does not need the runtime loop (`convey/voice.py:29-55`, `convey/voice.py:58-123`). + +### 3.2 Daily Briefing dispatch + +```text +cortex subprocess + -> think.cortex emits finish event to Callosum + -> push runtime listener callback receives message + -> triggers.handle_briefing_finish(message) + -> schedule coroutine on push runtime loop + -> poll _load_briefing_md(today) up to 10 x 1s + -> dispatch.send_many(eligible_devices, briefing_payload) + -> append journal/push/nudge_log.jsonl +``` + +The polling step is mandatory because `cortex.finish` is broadcast before `_write_output(...)` runs and before the `_active.jsonl` file is renamed to its completed name. `_load_briefing_md(today)` already enforces the `type=="morning_briefing"` and `metadata.date==today` gates, so the trigger reuses it instead of inventing a second briefing reader (`think/cortex.py:433-441`, `think/cortex.py:461-510`, `think/cortex.py:621-626`, `apps/home/routes.py:149-198`). + +### 3.3 Pre-Meeting Prep dispatch + +```text +push runtime periodic task (every 60s) + -> triggers.check_pre_meeting_prep(now) + -> get_enabled_facets().keys() + -> load_activity_records(facet, YYYYMMDD) + -> keep rows where source == "anticipated" + -> keep rows where start-now is within [14m, 16m] + -> skip rows already present in nudge_log + -> dispatch.send_many(eligible_devices, meeting_payload) + -> append journal/push/nudge_log.jsonl +``` + +The trigger must use `get_enabled_facets().keys()` so muted facets never produce push. The activity scan follows the existing repo convention: load all rows for `(facet, day)`, then filter `record["source"] == "anticipated"` in-process (`think/facets.py:255-261`, `think/activities.py:877-890`, `think/activities.py:937-939`, `apps/home/routes.py:305-337`). + +### 3.4 Agent Alert dispatch + +```text +in-process caller + -> triggers.send_agent_alert(title, body, context_id) + -> dispatch.send_many(eligible_devices, alert_payload) + -> append journal/push/nudge_log.jsonl +``` + +Agent Alert is intentionally the simplest path: no callosum subscription, no scheduler, no extra persistence besides dedupe log. It is a public in-process API for future callers that want to fire a push without adding new transport plumbing. + +### 3.5 Shutdown + +```text +process exit or test cleanup + -> atexit / stop_push_runtime(app) + -> stop CallosumConnection listener + -> cancel periodic asyncio task + -> stop runtime loop + -> join daemon thread + -> clear module-level runtime state +``` + +The shutdown contract matches the voice runtime: cancel tracked work, stop the loop from the owning thread, join with a bounded timeout, and leave the singleton reusable for the next app instance (`think/voice/runtime.py:53-109`, `tests/test_voice_runtime.py:19-103`). + +## 4. Endpoint specs + +All four endpoints inherit the default Convey auth gate because `convey/root.py` wraps every request unless the endpoint is on the explicit bypass allowlist, and new `/api/push/*` routes are not on that list (`convey/root.py:81-139`). Request validation mirrors the voice helpers in `convey/voice.py`, so malformed JSON and non-object bodies fail before feature code runs (`convey/voice.py:29-55`). + +### 4.1 `POST /api/push/register` + +- URL: `/api/push/register` +- Method: `POST` +- Auth: default gate only. Accepts a logged-in session, Basic Auth, or the existing `trust_localhost` bypass when setup is complete and proxy headers are absent (`convey/root.py:111-139`). +- Request schema: + - Required JSON object. + - `device_token: string` — trimmed, non-empty. The server strips embedded spaces and lowercases before storage. + - `bundle_id: string` — trimmed, non-empty. + - `environment: string` — must be `"development"` or `"production"`. + - `platform: string` — must be `"ios"` in Wave 3. +- Success response: + - HTTP 200 + - Body: `{"registered": true, "device_count": }` +- Error cases: + - HTTP 400 `{"error": "request body must be valid JSON"}` + - HTTP 400 `{"error": "request body must be a JSON object"}` + - HTTP 400 `{"error": "device_token is required"}` + - HTTP 400 `{"error": "bundle_id is required"}` + - HTTP 400 `{"error": "environment must be development or production"}` + - HTTP 400 `{"error": "platform must be ios"}` + - HTTP 500 `{"error": "device registration failed"}` +- Notes: + - Duplicate tokens upsert in place instead of adding a second row. + - `device_count` reports stored rows after the upsert. + - Registration does not require APNs credentials to be configured. Clients can register before the operator finishes `journal.json`. + +### 4.2 `DELETE /api/push/register` + +- URL: `/api/push/register` +- Method: `DELETE` +- Auth: default gate only (`convey/root.py:111-139`). +- Request schema: + - Required JSON object. + - `device_token: string` — trimmed, non-empty; normalized with the same rules as register. +- Success response: + - HTTP 200 + - Body: `{"removed": true, "device_count": }` when the token existed. + - HTTP 200 + - Body: `{"removed": false, "device_count": }` when the token was not present. +- Error cases: + - HTTP 400 `{"error": "request body must be valid JSON"}` + - HTTP 400 `{"error": "request body must be a JSON object"}` + - HTTP 400 `{"error": "device_token is required"}` + - HTTP 500 `{"error": "device removal failed"}` +- Notes: + - Removing an unknown token is not an error because uninstall and token churn are expected. + - The mutation rewrites `push_devices.json` only when the stored set actually changes. + +### 4.3 `GET /api/push/status` + +- URL: `/api/push/status` +- Method: `GET` +- Auth: default gate only (`convey/root.py:111-139`). +- Request schema: + - No request body. +- Success response: + - HTTP 200 + - Body: + - `configured: bool` + - `device_count: int` + - `devices: [{token_suffix, bundle_id, platform, environment, registered_at}]` +- `configured` semantics: + - `true` only when `push.apns_key_path`, `push.apns_key_id`, `push.apns_team_id`, and `push.bundle_id` are non-empty, `push.environment` resolves to `"development"` or `"production"`, and the configured `.p8` path is absolute, exists, and is readable. + - `false` for any missing or invalid APNs config, including a relative or missing key file. +- `devices` semantics: + - `token_suffix` is the last four characters of the stored token. + - `registered_at` is an ISO-8601 timestamp string. + - The list is sorted newest-first by `registered_at`. +- Error cases: + - No endpoint-specific 4xx errors. + - Malformed `push_devices.json` is treated as an empty store with a warning from `think/push/devices.py`; the route still returns HTTP 200. +- Notes: + - This route never returns full device tokens. + - This route reports all stored devices, not just devices matching the active bundle/environment filter. + +### 4.4 `POST /api/push/test` + +- URL: `/api/push/test` +- Method: `POST` +- Auth: default gate only. There is no extra debug header or test-only bypass (`convey/root.py:111-139`). +- Request schema: + - Empty body allowed. + - If a body is present, it must decode to a JSON object. + - Optional `title: string` + - Optional `body: string` + - Optional `category: string` +- Category rules: + - Default: `SOLSTONE_AGENT_ALERT` + - Allowed explicit values: `SOLSTONE_DAILY_BRIEFING`, `SOLSTONE_PRE_MEETING_PREP`, `SOLSTONE_AGENT_ALERT`, `SOLSTONE_COMMITMENT_NUDGE` + - The route validates the category value but still dispatches through `send_agent_alert(...)`, so Wave 3 test sends always use the Agent Alert payload shape. +- Success response: + - HTTP 200 + - Body: `{"sent": , "failed": }` +- Error cases: + - HTTP 400 `{"error": "request body must be valid JSON"}` + - HTTP 400 `{"error": "request body must be a JSON object"}` + - HTTP 400 `{"error": "category must be a known push category"}` + - HTTP 503 `{"error": "push not configured"}` +- Notes: + - The route dispatches through `send_agent_alert(...)`, so successful test sends append `SOLSTONE_AGENT_ALERT` lines to `nudge_log.jsonl` with unique `push-test-` context ids. + - The route sends only to stored devices whose `bundle_id`, `environment`, and `platform` match the current push configuration. + - Empty device store is a normal case and returns `{"sent": 0, "failed": 0}`. + +## 5. Config keys + +Push configuration lives only in `journal/config/journal.json`. Unlike voice, it does not fall back to environment variables, because these values describe the current journal’s push environment rather than a process-global secret cache (`think/voice/config.py:30-42`). Every read uses `get_config()` and every path lookup uses `get_journal()` at call time; the runtime does not cache the journal root during startup because `create_app()` does not assign `state.journal_root` until after `start_voice_runtime(app)` today, and Wave 3 preserves that ordering when it adds `start_push_runtime(app)` (`convey/__init__.py:163-166`, `think/voice/runtime.py:44-46`). + +| Key | Type | Default | Meaning | +|---|---|---|---| +| `push.apns_key_path` | `string \| null` | `null` | Absolute path to the APNs `.p8` signing key file. | +| `push.apns_key_id` | `string \| null` | `null` | Apple-issued APNs Key ID. | +| `push.apns_team_id` | `string \| null` | `null` | Apple Developer Team ID. | +| `push.bundle_id` | `string \| null` | `null` | App bundle identifier, for example `org.solpbc.solstone-swift`. | +| `push.environment` | `"development" \| "production" \| null` | `"development"` | APNs environment. `null` resolves to `"development"`. | + +Literal `think/journal_default.json` block: + +```json + "voice": { + "openai_api_key": null, + "model": "gpt-realtime", + "brain_model": "haiku" + }, + "push": { + "apns_key_path": null, + "apns_key_id": null, + "apns_team_id": null, + "bundle_id": null, + "environment": "development" + }, + "retention": { + "raw_media": "days", + "raw_media_days": 7, + "per_stream": {}, + "storage_warning_disk_percent": 80, + "storage_warning_raw_media_gb": null + } +``` + +Validation rules: + +- Blank strings normalize to `null`. +- `push.environment == null` resolves to `"development"`. +- Any non-null environment outside `{"development", "production"}` is invalid and makes dispatch unavailable until corrected. +- `push.apns_key_path` must be an absolute path. Relative paths are rejected because the server already has one journal-root concept and should not invent a second config root. +- `push.apns_key_path` must point to a readable file before the server reports push as configured. + +## 6. Payload shapes per category + +Common transport rules: + +- `dispatch.send(...)` uses `httpx.AsyncClient(http2=True)` against `https://api.sandbox.push.apple.com` for `"development"` and `https://api.push.apple.com` for `"production"`. +- Authorization is APNs bearer JWT signed with ES256 over the configured `.p8` key, using header `{"alg":"ES256","kid":}` and claims `{"iss":,"iat":}`. +- The JWT is cached in memory and regenerated when older than 55 minutes. Wave 3 does not mint a fresh token per request. +- `apns-topic` is always the configured `push.bundle_id`. +- `apns-priority` is `10` for every Wave 3 send because every shipped payload includes a visible alert. `5` remains reserved for future silent-only pushes and is not used in Wave 3. +- `BadDeviceToken` and `Unregistered` responses from APNs cause `dispatch.send(...)` to call `devices.remove_device(token)` before returning failure. + +Category constants defined by the server: + +- `SOLSTONE_DAILY_BRIEFING` +- `SOLSTONE_PRE_MEETING_PREP` +- `SOLSTONE_AGENT_ALERT` +- `SOLSTONE_COMMITMENT_NUDGE` + +### 6.1 Daily Briefing + +Headers: + +- `apns-topic: ` +- `apns-collapse-id: briefing.` +- `apns-priority: 10` + +Payload: + +```json +{ + "aps": { + "alert": { + "title": "Daily Briefing", + "body": "Your briefing is ready — tap to view" + }, + "category": "SOLSTONE_DAILY_BRIEFING", + "sound": "default", + "mutable-content": 1, + "content-available": 1 + }, + "data": { + "action": "open_briefing", + "day": "20260419", + "generated": "2026-04-19T06:45:00", + "needs_attention_count": 3 + } +} +``` + +Builder rules: + +- `day` is the local journal day in `YYYYMMDD` format. +- `generated` comes from briefing frontmatter when present, because `_load_briefing_md(today)` already loads that metadata (`apps/home/routes.py:149-198`, `tests/fixtures/journal/identity/briefing.md:1-14`). +- `needs_attention_count` is the length of the bullets list returned by `_load_briefing_md(today)` (`apps/home/routes.py:191-198`). + +### 6.2 Pre-Meeting Prep + +Headers: + +- `apns-topic: ` +- `apns-collapse-id: meeting.` +- `apns-priority: 10` + +Payload: + +```json +{ + "aps": { + "alert": { + "title": "Pre-Meeting Prep", + "body": "Meeting in 15 minutes — tap to view" + }, + "category": "SOLSTONE_PRE_MEETING_PREP", + "sound": "default", + "mutable-content": 1, + "content-available": 1, + "interruption-level": "time-sensitive" + }, + "data": { + "action": "open_pre_meeting", + "activity_id": "anticipated_meeting_090000_0420", + "facet": "work", + "day": "20260420", + "start": "09:00", + "title": "Launch sync", + "location": "Room A", + "participants": [ + "Juliet Capulet" + ], + "prep_notes": "Bring launch notes" + } +} +``` + +Builder rules: + +- `activity_id` comes from the activity record’s `id`, which already exists on anticipated rows and is the right natural key for collapse and dedupe (`tests/test_voice_tools.py:197-211`, `think/activities.py:893-915`). +- `start` accepts stored `HH:MM` or `HH:MM:SS`; the trigger parser is tolerant, but the payload preserves the stored string as-is. +- `participants` pulls attendee names from `participation` entries where `role=="attendee"`, following the existing Home surface pattern (`apps/home/routes.py:312-337`, `tests/test_voice_tools.py:197-214`). +- `interruption-level` is present only on this category. + +### 6.3 Agent Alert + +Headers: + +- `apns-topic: ` +- `apns-collapse-id: alert.` +- `apns-priority: 10` + +Payload: + +```json +{ + "aps": { + "alert": { + "title": "Agent Alert", + "body": "A workflow needs attention" + }, + "category": "SOLSTONE_AGENT_ALERT", + "sound": "default", + "mutable-content": 1, + "content-available": 1 + }, + "data": { + "action": "open_alert", + "context_id": "triage-20260419-001" + } +} +``` + +Builder rules: + +- `title` and `body` come from the caller. +- `context_id` is required for the public `send_agent_alert(...)` helper and is used directly in both `data.context_id` and the collapse id. +- `POST /api/push/test` generates a UUID context when the client does not supply one through the route body. + +### 6.4 Deferred Commitment Nudge + +Headers: + +- `apns-topic: ` +- `apns-collapse-id: commitment.` +- `apns-priority: 10` + +Forward-compatible payload: + +```json +{ + "aps": { + "alert": { + "title": "Commitment Nudge", + "body": "A commitment needs attention — tap to view" + }, + "category": "SOLSTONE_COMMITMENT_NUDGE", + "sound": "default", + "mutable-content": 1, + "content-available": 1 + }, + "data": { + "action": "open_commitment", + "ledger_id": "lg_123" + } +} +``` + +Wave 3 defines this payload shape and constant only for client forward-compatibility. No trigger emits it until Wave 3.1 lands a real due-date primitive. + +### 6.5 PII fallback rule + +Wave 3 hardcodes the lock-screen-safe fallback into the payload builders. It is not a TODO for the iOS client. + +- Daily Briefing body is always generic: “Your briefing is ready — tap to view.” +- Pre-Meeting Prep body is always generic: “Meeting in 15 minutes — tap to view.” +- Deferred Commitment Nudge body is always generic: “A commitment needs attention — tap to view.” +- Sensitive detail lives in `data`, where the Notification Service Extension can read it off-device and decide what to reveal. +- Agent Alert is the exception: `title` and `body` are caller-provided, so the caller is responsible for keeping them lock-screen safe. + +## 7. Domain write-ownership (L1–L9 declarations) + +Push owns its own journal domain. It is not an indexer, importer, scheduler, or search subsystem, so its writes belong in `think/push/*` and do not need to be routed through another domain owner (`scripts/check_layer_hygiene.py:38-57`, `scripts/check_layer_hygiene.py:199-209`). + +| Path | Owner module | Write API | Read API | +|---|---|---|---| +| `journal/config/push_devices.json` | `think/push/devices.py` | `register_device`, `remove_device` | `load_devices` | +| `journal/push/nudge_log.jsonl` | `think/push/triggers.py` | `_append_nudge_log` | `_has_nudged` | + +L1 declaration: + +- `think/push/*` is a feature-owned runtime and transport layer, not infrastructure. Its journal writes are feature state, not accidental cross-domain mutation. + +L2 declaration: + +- No module outside `think/push/` writes `journal/config/push_devices.json`. +- No module outside `think/push/` writes `journal/push/nudge_log.jsonl`. +- `convey/push.py` calls feature APIs but never writes these paths directly. + +L3 declaration: + +- `load_devices()` never writes. Malformed-store recovery is explicit: the reader returns an empty list and leaves rewrite responsibility to the next write path. +- `_has_nudged(...)` never writes. Dedupe is read-first, then write on success. +- There is no create-on-miss hidden behind any `load_*` or `get_*` helper. + +L6/L7 declaration: + +- Indexers and importers do not touch push paths. +- The current hygiene script only scans infrastructure scopes `think/indexer`, `think/importers`, `think/search`, and `think/graph`, plus read-verb `apps/*/call.py` handlers; `think/push/*` is outside those scopes (`scripts/check_layer_hygiene.py:38-44`, `scripts/check_layer_hygiene.py:124-145`, `scripts/check_layer_hygiene.py:156-180`). +- The current hygiene script also only looks for writes near `journal/entities`, `journal/facets`, and `journal/observations`, not `journal/config/push_devices.json` or `journal/push/nudge_log.jsonl`, so no allowlist entry is required (`scripts/check_layer_hygiene.py:59-72`, `scripts/check_layer_hygiene.py:105-108`). + +L8 declaration: + +- Hooks do not apply. Push has no talent hook and does not write through `think/hooks.py` or `talent/*.py`. + +L9 declaration: + +- Daily Briefing dedupe key: `(SOLSTONE_DAILY_BRIEFING, YYYYMMDD)` +- Pre-Meeting Prep dedupe key: `(SOLSTONE_PRE_MEETING_PREP, activity_id, YYYYMMDD)` +- Agent Alert dedupe key: `(SOLSTONE_AGENT_ALERT, context_id)` +- Commitment Nudge reserved dedupe key: `(SOLSTONE_COMMITMENT_NUDGE, ledger_id)` +- Trigger order is always: + - compute dedupe key + - check `_has_nudged(...)` + - send to eligible devices + - append log only when `sent > 0` +- That write-after-success rule keeps Pre-Meeting Prep retryable if the first attempt produces zero successful sends, while still preventing duplicate delivery after one device has already received the notification. + +`nudge_log.jsonl` line shape: + +- Daily Briefing line: `ts`, `category`, `dedupe_key`, `day`, `sent`, `failed` +- Pre-Meeting Prep line: `ts`, `category`, `dedupe_key`, `day`, `facet`, `activity_id`, `sent`, `failed` +- Agent Alert line: `ts`, `category`, `dedupe_key`, `context_id`, `sent`, `failed` + +## 8. Tests + +All push tests use `tests/fixtures/journal/` plus `monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", ...)` where necessary, following the existing voice integration and route test setup pattern (`tests/test_voice_routes.py:14-23`, `tests/test_voice_integration.py:102-149`). + +### 8.1 `tests/test_push_config.py` + +Purpose: mirror the small-reader coverage shape of `tests/test_voice_config.py` while locking down the journal-only push config contract (`tests/test_voice_config.py:9-43`). + +- Defaults: all four required APNs fields resolve to `None`, and `push.environment` resolves to `"development"` when omitted. +- Whitespace cleanup: blank strings normalize to `None`. +- Journal precedence: populated `journal.json` values are used directly. +- No env fallback: setting unrelated env vars does not change push config reads. +- Invalid environment: non-`development` / non-`production` values raise a validation error or mark config unavailable, depending on the caller. +- Missing key path: status/configured becomes false when the configured `.p8` path does not exist. + +### 8.2 `tests/test_push_devices.py` + +Purpose: cover the device store as its own write-owning module. + +- Register/load/remove round trip against a temporary journal. +- Duplicate token registration updates the existing row instead of creating a second device. +- `registered_at` refreshes on duplicate registration. +- Remove returns `False` for unknown token and leaves the count unchanged. +- Empty store returns `[]` and does not create the file. +- Malformed `push_devices.json` is treated as empty with a warning, then repaired by the next successful write. +- Status-shaping helper masks tokens to last four only. + +### 8.3 `tests/test_push_dispatch.py` + +Purpose: unit-test APNs transport, JWT handling, payload builders, and token-redaction rules. + +- JWT signing shape: header includes `alg=ES256` and `kid=`, claims include `iss=` and `iat=`. +- JWT caching: second send inside 55 minutes reuses the cached token. +- JWT refresh: send at 60+ minutes mints a fresh token. +- Payload shape for Daily Briefing matches §6 exactly. +- Payload shape for Pre-Meeting Prep includes `interruption-level: time-sensitive`. +- Payload shape for Agent Alert matches §6 exactly. +- Commitment payload constant and builder exist but are not wired to a trigger. +- `apns-collapse-id` is `briefing.`, `meeting.`, `alert.`, and `commitment.`. +- `BadDeviceToken` response calls `devices.remove_device(token)`. +- `Unregistered` response calls `devices.remove_device(token)`. +- Captured logs never include full tokens, JWTs, or raw `.p8` contents. + +### 8.4 `tests/test_push_triggers.py` + +Purpose: lock down idempotency, timing, and source filtering in the trigger layer. + +- Daily Briefing finish polls until `briefing.md` exists, using a delayed write to simulate the `cortex.finish` before-write ordering. +- Daily Briefing gives up after 10 polls and logs a warning. +- Daily Briefing fires once per day even if the same finish message is delivered twice. +- Daily Briefing uses `_load_briefing_md(today)` and ignores stale or wrong-type briefing files. +- Pre-Meeting Prep scans only `get_enabled_facets().keys()` and skips muted facets. +- Pre-Meeting Prep filters `record["source"] == "anticipated"` and ignores all other activity rows. +- Pre-Meeting Prep detects the 15-minute window with both `HH:MM` and `HH:MM:SS` starts. +- Pre-Meeting Prep is idempotent across repeated periodic ticks within the same 2-minute window. +- `send_agent_alert(...)` builds the expected payload and appends the expected dedupe log line. +- Zero-device and zero-success sends do not append dedupe log lines. + +### 8.5 `tests/test_push_routes.py` + +Purpose: mirror the validation-heavy style of `tests/test_voice_routes.py` at the new `/api/push/*` surface (`tests/test_voice_routes.py:26-118`). + +- `POST /api/push/register` happy path. +- `POST /api/push/register` rejects non-object JSON. +- `POST /api/push/register` rejects missing fields and invalid `environment` / `platform`. +- `DELETE /api/push/register` happy path for an existing token. +- `DELETE /api/push/register` returns `removed: false` for a missing token. +- `GET /api/push/status` returns `configured`, `device_count`, and masked device rows. +- `POST /api/push/test` returns `503` when APNs config is missing. +- `POST /api/push/test` aggregates sent/failed counts from the dispatch layer. +- `POST /api/push/test` rejects unknown category strings. + +### 8.6 `tests/test_push_integration.py` + +Purpose: prove the whole runtime path from Flask startup through callosum listener and periodic dispatch, following the same end-to-end philosophy as `tests/test_voice_integration.py` (`tests/test_voice_integration.py:115-149`). + +- Boot `create_app()` against a fixture journal and verify `start_push_runtime(app)` runs. +- Patch `httpx.AsyncClient.post` so APNs sends are fully mocked. +- Patch JWT signing inputs so tests can assert stable Authorization headers. +- Fire a fake `cortex.finish` event through the push runtime’s `CallosumConnection` callback and assert Daily Briefing dispatch occurs. +- Confirm `nudge_log.jsonl` receives the expected Daily Briefing line. +- Seed anticipated activities and run one periodic `check_pre_meeting_prep(now)` pass through the runtime loop. +- Confirm muted facets do not dispatch. +- Stop the runtime cleanly and assert the loop and thread are cleared. + +## 9. Security considerations + +- Device-token redaction: full tokens are never returned by `GET /api/push/status` and never written to logs. Status exposes last four characters only, and dispatch logs use the same redaction rule. +- APNs JWT secrecy: bearer JWT values are never logged. Refresh decisions may log age or cache-hit state, but not the token string itself. +- `.p8` key secrecy: `push.apns_key_path` may appear in operator config, but the file contents themselves are never logged or echoed back through routes. +- PII fallback: Daily Briefing, Pre-Meeting Prep, and deferred Commitment Nudge use generic lock-screen bodies, with detail moved into `data` for device-side handling. Agent Alert is caller-provided and inherits caller responsibility for lock-screen safety. +- Auth: all `/api/push/*` endpoints use the default root auth gate. There is no debug header and no push-specific bypass. Default auth is session cookie, then Basic Auth, then opt-in `trust_localhost` after setup when proxy headers are absent (`convey/root.py:49-57`, `convey/root.py:81-139`). +- `trust_localhost` stays narrow by design: it only applies after setup completion and only when `request.remote_addr` is local and proxy headers are absent (`convey/root.py:119-139`). +- Facet eligibility: all push triggers operate on `get_enabled_facets().keys()`, so muted facets are excluded from both dispatch and device-visible summaries (`think/facets.py:255-261`, `think/surfaces/ledger.py:454-456`). +- Terminology covenant: operator-visible strings and payload labels use the repo’s “observer/listen” vocabulary and avoid “capture”, “record”, “keeper”, or “assistant”. +- Hosted-MVP privacy stance: payloads are cleartext to APNs and the device; Wave 3 is explicitly non-E2E. +- No analytics: Wave 3 adds no tracking, analytics beacons, crash reporting, or delivery pixel equivalents. + +## 10. Live validation + +**Live APNs validation is deferred pending Apple Developer enrollment.** Wave 3 ships infrastructure and mocked tests only. When enrollment completes: + +1. Configure `push.apns_key_path`, `push.apns_key_id`, `push.apns_team_id`, `push.bundle_id`, and `push.environment` in `journal.json`. +2. Register a real device from the iOS client. +3. Exercise `POST /api/push/test` against the development APNs environment and confirm the device receives the notification. +4. Manually trigger a `morning_briefing` cortex run and confirm the Daily Briefing push lands. +5. Wait for a scheduled meeting and confirm Pre-Meeting Prep lands within approximately ±30 seconds of T-15:00. +6. Only then flip `push.environment` to `production` and deploy. + +Sandbox smoke-test commands: + +```sh +BASE_URL=${BASE_URL:-http://127.0.0.1:5015} +AUTH=${AUTH:-":$SOL_PASSWORD"} +TOKEN=${TOKEN:-0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef} + +curl -u "$AUTH" \ + -H 'Content-Type: application/json' \ + -X POST "$BASE_URL/api/push/register" \ + -d '{ + "device_token": "'"$TOKEN"'", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios" + }' + +curl -u "$AUTH" \ + "$BASE_URL/api/push/status" + +curl -u "$AUTH" \ + -H 'Content-Type: application/json' \ + -X POST "$BASE_URL/api/push/test" \ + -d '{ + "title": "Push test", + "body": "This is a sandbox test notification.", + "category": "SOLSTONE_AGENT_ALERT" + }' + +curl -u "$AUTH" \ + -H 'Content-Type: application/json' \ + -X DELETE "$BASE_URL/api/push/register" \ + -d '{ + "device_token": "'"$TOKEN"'" + }' +``` + +Basic Auth uses only the password component, so `-u ":$SOL_PASSWORD"` is the portable curl form for these routes (`convey/root.py:49-57`). + +## 11. Open questions + +- Agent Alert body limits: Wave 3 should probably enforce a soft cap before the native client ships, but the exact truncation policy can wait until the iOS notification UI settles. +- Multi-build device identity: Wave 3 keys stored devices by token alone. If one journal starts registering multiple app builds that share a token namespace, revisit whether identity should widen to `(token, bundle_id, environment)`. +- Retry telemetry: Wave 3 records dedupe state in `nudge_log.jsonl`, but it does not yet record APNs failure reasons in a separate operator-facing history file. + +## 12. Sources + +Voice-server analog: + +- `docs/design/voice-server.md:1-465`, `convey/voice.py:26-184`, `convey/__init__.py:112-166`, `think/voice/runtime.py:21-109`, `think/voice/config.py:17-42`, `think/voice/sideband.py:20-61`, `tests/test_voice_config.py:9-43`, `tests/test_voice_runtime.py:19-103`, `tests/test_voice_routes.py:26-118`, `tests/test_voice_integration.py:102-149` + +Callosum / cortex finish timing: + +- `think/callosum.py:245-346`, `think/cortex.py:433-441`, `think/cortex.py:461-510`, `think/cortex.py:621-626`, `apps/home/routes.py:149-198`, `apps/home/workspace.html:1468-1470`, `apps/home/workspace.html:1732-1788`, `convey/bridge.py:45-86`, `apps/home/events.py:21-55` + +Ledger / activities APIs: + +- `think/surfaces/types.py:16-31`, `think/surfaces/ledger.py:395-487`, `think/activities.py:877-945`, `apps/home/routes.py:305-337`, `think/facets.py:255-261`, `tests/test_voice_tools.py:197-214`, `tests/fixtures/journal/identity/briefing.md:1-14` + +Auth model: + +- `convey/root.py:49-57`, `convey/root.py:81-139` + +Scheduler / runtime choice: + +- `think/scheduler.py:29-30`, `think/scheduler.py:375-438`, `think/supervisor.py:1311-1371`, `think/heartbeat.py:45-138` + +Layer-hygiene script: + +- `scripts/check_layer_hygiene.py:38-72`, `scripts/check_layer_hygiene.py:105-108`, `scripts/check_layer_hygiene.py:124-145`, `scripts/check_layer_hygiene.py:156-180`, `scripts/check_layer_hygiene.py:199-220` + +Config and dependencies: + +- `think/journal_default.json:35-39`, `pyproject.toml:53` + +Wave 2 commit: + +- `0a693381 voice: ship Wave 2 voice server (root /api/voice/*, 9-tool sideband)` diff --git a/pyproject.toml b/pyproject.toml index ae86343f2..e6c180a69 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -51,6 +51,8 @@ dependencies = [ "openai-agents>=0.1.0", "anthropic", "httpx", + "h2", + "pyjwt>=2.8", "jsonschema>=4.26,<5", "genai-prices", # Link tunnel service (think/link/): TLS 1.3 in memory-BIO mode over diff --git a/tests/conftest.py b/tests/conftest.py index 36603bfe6..dc9d723bf 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -19,6 +19,7 @@ from think.entities.journal import clear_journal_entity_cache from think.entities.loading import clear_entity_loading_cache from think.entities.observations import clear_observation_cache from think.entities.relationships import clear_relationship_caches +from think.push.runtime import stop_all_push_runtime from think.utils import now_ms from think.voice import brain as voice_brain from think.voice.runtime import stop_all_voice_runtime @@ -67,6 +68,12 @@ def _cleanup_voice_runtime(): voice_brain.clear_brain_state() +@pytest.fixture(autouse=True) +def _cleanup_push_runtime(): + yield + stop_all_push_runtime() + + @pytest.fixture def journal_copy(tmp_path, monkeypatch): """Copy git-tracked fixture files to tmp_path for mutation tests.""" diff --git a/tests/test_push_config.py b/tests/test_push_config.py new file mode 100644 index 000000000..2e05b2ce2 --- /dev/null +++ b/tests/test_push_config.py @@ -0,0 +1,140 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from think.push import config + + +def _write_config(tmp_path: Path, payload: dict) -> None: + config_path = tmp_path / "config" / "journal.json" + config_path.parent.mkdir(parents=True, exist_ok=True) + config_path.write_text(json.dumps(payload), encoding="utf-8") + + +def test_push_config_defaults(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + _write_config(tmp_path, {"agent": {"name": "sol"}}) + + assert config.get_apns_key_path() is None + assert config.get_apns_key_id() is None + assert config.get_apns_team_id() is None + assert config.get_bundle_id() is None + assert config.get_environment() == "development" + assert config.is_configured() is False + + +def test_push_config_reads_journal_values(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + key_path = tmp_path / "keys" / "apns.p8" + key_path.parent.mkdir(parents=True, exist_ok=True) + key_path.write_text("PRIVATE KEY", encoding="utf-8") + _write_config( + tmp_path, + { + "push": { + "apns_key_path": f" {key_path} ", + "apns_key_id": " KEY123 ", + "apns_team_id": " TEAM123 ", + "bundle_id": " org.solpbc.solstone-swift ", + "environment": "production", + } + }, + ) + + assert config.get_apns_key_path() == key_path + assert config.get_apns_key_id() == "KEY123" + assert config.get_apns_team_id() == "TEAM123" + assert config.get_bundle_id() == "org.solpbc.solstone-swift" + assert config.get_environment() == "production" + assert config.is_configured() is True + + +def test_push_config_blank_values_normalize_to_none(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + _write_config( + tmp_path, + { + "push": { + "apns_key_path": " ", + "apns_key_id": "\t", + "apns_team_id": "", + "bundle_id": " ", + "environment": " ", + } + }, + ) + + assert config.get_apns_key_path() is None + assert config.get_apns_key_id() is None + assert config.get_apns_team_id() is None + assert config.get_bundle_id() is None + assert config.get_environment() == "development" + assert config.is_configured() is False + + +def test_push_config_invalid_environment_raises(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + _write_config(tmp_path, {"push": {"environment": "staging"}}) + + with pytest.raises( + ValueError, match="push.environment must be 'development' or 'production'" + ): + config.get_environment() + + assert config.is_configured() is False + + +def test_push_config_missing_key_file_is_unconfigured(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + _write_config( + tmp_path, + { + "push": { + "apns_key_path": str(tmp_path / "missing.p8"), + "apns_key_id": "KEY123", + "apns_team_id": "TEAM123", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + } + }, + ) + + assert config.is_configured() is False + + +def test_push_config_relative_key_path_is_unconfigured(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + relative_key_path = Path("keys/apns.p8") + _write_config( + tmp_path, + { + "push": { + "apns_key_path": str(relative_key_path), + "apns_key_id": "KEY123", + "apns_team_id": "TEAM123", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + } + }, + ) + + assert config.get_apns_key_path() == relative_key_path + assert config.is_configured() is False + + +def test_push_config_ignores_env_fallback(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + monkeypatch.setenv("OPENAI_API_KEY", "sk-env") + monkeypatch.setenv("APNS_KEY_ID", "ENVKEY") + _write_config(tmp_path, {"push": {}}) + + assert config.get_apns_key_id() is None + assert config.get_apns_team_id() is None + assert config.get_bundle_id() is None + assert config.is_configured() is False diff --git a/tests/test_push_devices.py b/tests/test_push_devices.py new file mode 100644 index 000000000..478a22935 --- /dev/null +++ b/tests/test_push_devices.py @@ -0,0 +1,136 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +from pathlib import Path + +from think.push import devices + + +def _devices_path(tmp_path: Path) -> Path: + return tmp_path / "config" / "push_devices.json" + + +def test_load_devices_returns_empty_for_missing_store(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + + assert devices.load_devices() == [] + + +def test_register_load_remove_round_trip(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + + count = devices.register_device( + token="a" * 64, + bundle_id="org.solpbc.solstone-swift", + environment="development", + platform="ios", + ) + + assert count == 1 + stored = devices.load_devices() + assert stored == [ + { + "token": "a" * 64, + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios", + "registered_at": stored[0]["registered_at"], + } + ] + + removed = devices.remove_device("a" * 64) + assert removed is True + assert devices.load_devices() == [] + + +def test_register_device_updates_existing_token(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + times = iter([1000, 2000]) + monkeypatch.setattr(devices.time, "time", lambda: next(times)) + + first = devices.register_device( + token="b" * 64, + bundle_id="org.solpbc.solstone-swift", + environment="development", + platform="ios", + ) + second = devices.register_device( + token="b" * 64, + bundle_id="org.solpbc.solstone-swift", + environment="production", + platform="ios", + ) + + assert first == 1 + assert second == 1 + assert devices.load_devices() == [ + { + "token": "b" * 64, + "bundle_id": "org.solpbc.solstone-swift", + "environment": "production", + "platform": "ios", + "registered_at": 2000, + } + ] + + +def test_remove_device_returns_false_for_unknown_token(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + devices.register_device( + token="c" * 64, + bundle_id="org.solpbc.solstone-swift", + environment="development", + platform="ios", + ) + + assert devices.remove_device("d" * 64) is False + assert len(devices.load_devices()) == 1 + + +def test_load_devices_returns_empty_for_malformed_store(monkeypatch, tmp_path, caplog): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + path = _devices_path(tmp_path) + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text('{"devices": "bad"}', encoding="utf-8") + + loaded = devices.load_devices() + + assert loaded == [] + assert "push device store unreadable" in caplog.text + + +def test_status_view_masks_token(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + path = _devices_path(tmp_path) + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text( + json.dumps( + { + "devices": [ + { + "token": "0123456789abcdef", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios", + "registered_at": 1713528000, + } + ] + } + ), + encoding="utf-8", + ) + + device = devices.load_devices()[0] + view = devices.status_view(device) + + assert devices.mask_token("0123456789abcdef") == "...cdef" + assert view == { + "token_suffix": "...cdef", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios", + "registered_at": "2024-04-19T12:00:00Z", + } diff --git a/tests/test_push_dispatch.py b/tests/test_push_dispatch.py new file mode 100644 index 000000000..00751e5bf --- /dev/null +++ b/tests/test_push_dispatch.py @@ -0,0 +1,222 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import re +from pathlib import Path +from unittest.mock import patch + +import httpx +import jwt + +from think.push import dispatch + +TEST_KEY = """-----BEGIN PRIVATE KEY----- +MIGHAgEAMBMGByqGSM49AgEGCCqGSM49AwEHBG0wawIBAQQg+Zj7Bk6Dzp080/PU +jTZnJ6kP4KtlHErFO/WuVRTQvkShRANCAARW8djY5HF7K8noSZQRfjP38mIzaufi +/YPI38YuaWmiPIqRmwDOu5rICl4PPLem4k+qtb950rlYCGx3J+MQN9tO +-----END PRIVATE KEY----- +""" + + +def _write_key(tmp_path: Path) -> Path: + key_path = tmp_path / "keys" / "apns.p8" + key_path.parent.mkdir(parents=True, exist_ok=True) + key_path.write_text(TEST_KEY, encoding="utf-8") + return key_path + + +def _configure_push(monkeypatch, tmp_path: Path) -> None: + key_path = _write_key(tmp_path) + monkeypatch.setattr(dispatch, "get_apns_key_path", lambda: key_path) + monkeypatch.setattr(dispatch, "get_apns_key_id", lambda: "KEY123") + monkeypatch.setattr(dispatch, "get_apns_team_id", lambda: "TEAM123") + monkeypatch.setattr(dispatch, "get_bundle_id", lambda: "org.solpbc.solstone-swift") + monkeypatch.setattr(dispatch, "get_environment", lambda: "development") + dispatch._APNS_JWT_CACHE.clear() + + +def test_mint_apns_jwt_has_expected_header_and_claims(monkeypatch, tmp_path): + _configure_push(monkeypatch, tmp_path) + + token = dispatch._mint_apns_jwt(now=1713528000) + + assert jwt.get_unverified_header(token) == { + "alg": "ES256", + "kid": "KEY123", + "typ": "JWT", + } + assert jwt.decode(token, options={"verify_signature": False}) == { + "iss": "TEAM123", + "iat": 1713528000, + } + + +def test_mint_apns_jwt_reuses_cached_token_within_55_minutes(monkeypatch, tmp_path): + _configure_push(monkeypatch, tmp_path) + + first = dispatch._mint_apns_jwt(now=1000) + second = dispatch._mint_apns_jwt(now=1000 + 55 * 60) + + assert first == second + + +def test_mint_apns_jwt_refreshes_after_55_minutes(monkeypatch, tmp_path): + _configure_push(monkeypatch, tmp_path) + + first = dispatch._mint_apns_jwt(now=1000) + second = dispatch._mint_apns_jwt(now=1000 + 55 * 60 + 1) + + assert first != second + + +def test_daily_briefing_payload_shape(): + payload = dispatch.build_daily_briefing_payload( + day="20260419", generated="2026-04-19T06:45:00", needs_attention_count=3 + ) + + assert payload["aps"]["category"] == dispatch.CATEGORY_DAILY_BRIEFING + assert payload["aps"]["sound"] == "default" + assert payload["aps"]["mutable-content"] == 1 + assert payload["aps"]["content-available"] == 1 + assert "interruption-level" not in payload["aps"] + assert payload["data"] == { + "action": "open_briefing", + "day": "20260419", + "generated": "2026-04-19T06:45:00", + "needs_attention_count": 3, + } + + +def test_pre_meeting_payload_shape(): + payload = dispatch.build_pre_meeting_payload( + activity={ + "id": "anticipated_meeting_090000_0420", + "start": "09:00", + "title": "Launch sync", + "location": "Room A", + "prep_notes": "Bring launch notes", + "participation": [ + {"name": "Juliet Capulet", "role": "attendee"}, + {"name": "Observer", "role": "organizer"}, + ], + }, + facet="work", + day="20260420", + ) + + assert payload["aps"]["category"] == dispatch.CATEGORY_PRE_MEETING_PREP + assert payload["aps"]["interruption-level"] == "time-sensitive" + assert payload["data"]["action"] == "open_pre_meeting" + assert payload["data"]["participants"] == ["Juliet Capulet"] + + +def test_agent_alert_payload_shape(): + payload = dispatch.build_agent_alert_payload( + title="Agent Alert", body="Needs review", context_id="ctx-1" + ) + + assert payload["aps"]["category"] == dispatch.CATEGORY_AGENT_ALERT + assert payload["data"] == {"action": "open_alert", "context_id": "ctx-1"} + assert "interruption-level" not in payload["aps"] + + +def test_commitment_payload_shape(): + payload = dispatch.build_commitment_payload(ledger_id="lg_123") + + assert payload["aps"]["category"] == dispatch.CATEGORY_COMMITMENT_NUDGE + assert payload["data"] == {"action": "open_commitment", "ledger_id": "lg_123"} + + +def test_collapse_ids(): + assert dispatch.build_daily_briefing_collapse_id("20260419") == "briefing.20260419" + assert ( + dispatch.build_pre_meeting_collapse_id("anticipated_meeting_090000_0420") + == "meeting.anticipated_meeting_090000_0420" + ) + assert dispatch.build_agent_alert_collapse_id("ctx-1") == "alert.ctx-1" + assert dispatch.build_commitment_collapse_id("lg_123") == "commitment.lg_123" + + +def test_send_removes_bad_device_token(monkeypatch, tmp_path): + _configure_push(monkeypatch, tmp_path) + removed: list[str] = [] + monkeypatch.setattr( + dispatch.devices, "remove_device", lambda token: removed.append(token) or True + ) + + async def fake_post(self, url, *, headers, json): + return httpx.Response(400, json={"reason": "BadDeviceToken"}) + + with patch.object(httpx.AsyncClient, "post", new=fake_post): + ok, reason = dispatch.send( + {"token": "a" * 64}, + dispatch.build_agent_alert_payload( + title="Agent Alert", body="Needs review", context_id="ctx-1" + ), + collapse_id="alert.ctx-1", + ) + + assert ok is False + assert reason == "BadDeviceToken" + assert removed == ["a" * 64] + + +def test_send_removes_unregistered_device_on_410(monkeypatch, tmp_path): + _configure_push(monkeypatch, tmp_path) + removed: list[str] = [] + monkeypatch.setattr( + dispatch.devices, "remove_device", lambda token: removed.append(token) or True + ) + + async def fake_post(self, url, *, headers, json): + return httpx.Response(410, json={"reason": "Unregistered"}) + + with patch.object(httpx.AsyncClient, "post", new=fake_post): + ok, reason = dispatch.send( + {"token": "b" * 64}, + dispatch.build_agent_alert_payload( + title="Agent Alert", body="Needs review", context_id="ctx-1" + ), + collapse_id="alert.ctx-1", + ) + + assert ok is False + assert reason == "Unregistered" + assert removed == ["b" * 64] + + +def test_send_many_reuses_client_and_redacts_tokens(monkeypatch, tmp_path, caplog): + _configure_push(monkeypatch, tmp_path) + calls: list[dict[str, object]] = [] + caplog.set_level("WARNING", logger="solstone.push.dispatch") + + async def fake_post(self, url, *, headers, json): + calls.append({"url": url, "headers": headers, "json": json}) + return httpx.Response(500, json={"reason": "InternalServerError"}) + + with patch.object(httpx.AsyncClient, "post", new=fake_post): + sent, failed = dispatch.send_many( + [ + {"token": "c" * 64}, + {"token": "d" * 64}, + ], + dispatch.build_daily_briefing_payload( + day="20260419", + generated="2026-04-19T06:45:00", + needs_attention_count=1, + ), + collapse_id="briefing.20260419", + ) + + assert sent == 0 + assert failed == 2 + assert len(calls) == 2 + assert calls[0]["headers"]["apns-collapse-id"] == "briefing.20260419" + assert calls[0]["headers"]["apns-priority"] == "10" + assert calls[0]["headers"]["apns-push-type"] == "alert" + assert calls[0]["headers"]["apns-topic"] == "org.solpbc.solstone-swift" + assert "push rejected token=...cccc" in caplog.text + assert all(record.levelname == "WARNING" for record in caplog.records) + assert re.search(r"[0-9a-f]{64}", caplog.text) is None diff --git a/tests/test_push_integration.py b/tests/test_push_integration.py new file mode 100644 index 000000000..39474face --- /dev/null +++ b/tests/test_push_integration.py @@ -0,0 +1,146 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +from datetime import datetime +from pathlib import Path +from unittest.mock import patch + +import httpx + +from convey import create_app +from think.push import devices, runtime, triggers + +TEST_KEY = """-----BEGIN PRIVATE KEY----- +MIGHAgEAMBMGByqGSM49AgEGCCqGSM49AwEHBG0wawIBAQQg+Zj7Bk6Dzp080/PU +jTZnJ6kP4KtlHErFO/WuVRTQvkShRANCAARW8djY5HF7K8noSZQRfjP38mIzaufi +/YPI38YuaWmiPIqRmwDOu5rICl4PPLem4k+qtb950rlYCGx3J+MQN9tO +-----END PRIVATE KEY----- +""" + + +class FixedDateTime(datetime): + @classmethod + def now(cls, tz=None): + return cls(2026, 3, 27, 8, 45, 0, tzinfo=tz) + + +def _write_push_config(journal_copy: Path) -> None: + key_path = journal_copy / "keys" / "apns.p8" + key_path.parent.mkdir(parents=True, exist_ok=True) + key_path.write_text(TEST_KEY, encoding="utf-8") + config_path = journal_copy / "config" / "journal.json" + config = json.loads(config_path.read_text(encoding="utf-8")) + config["push"] = { + "apns_key_path": str(key_path), + "apns_key_id": "KEY123", + "apns_team_id": "TEAM123", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + } + config_path.write_text(json.dumps(config, indent=2) + "\n", encoding="utf-8") + + +def _seed_activity(journal_copy: Path, facet: str, day: str, rows: list[dict]) -> None: + path = journal_copy / "facets" / facet / "activities" / f"{day}.jsonl" + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text("\n".join(json.dumps(row) for row in rows) + "\n", encoding="utf-8") + + +def test_push_integration_briefing_dispatch_and_log(journal_copy, monkeypatch): + _write_push_config(journal_copy) + devices.register_device( + token="a" * 64, + bundle_id="org.solpbc.solstone-swift", + environment="development", + platform="ios", + ) + monkeypatch.setattr( + "think.push.runtime.CallosumConnection.start", + lambda self, callback=None: None, + ) + monkeypatch.setattr("think.push.runtime.CallosumConnection.stop", lambda self: None) + monkeypatch.setattr(triggers, "datetime", FixedDateTime) + captured: list[str] = [] + + async def fake_post(self, url, *, headers, json): + captured.append(url) + return httpx.Response(200) + + with patch.object(httpx.AsyncClient, "post", new=fake_post): + app = create_app(str(journal_copy)) + app.config["TESTING"] = True + runtime._on_callosum_message( + {"tract": "cortex", "event": "finish", "name": "morning_briefing"} + ) + + log_path = journal_copy / "push" / "nudge_log.jsonl" + lines = [ + json.loads(line) + for line in log_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + assert captured == [f"https://api.sandbox.push.apple.com/3/device/{'a' * 64}"] + assert lines[0]["category"] == "SOLSTONE_DAILY_BRIEFING" + runtime.stop_all_push_runtime() + assert runtime.get_runtime_state() is None + + +def test_push_integration_pre_meeting_and_muted_facet(journal_copy): + _write_push_config(journal_copy) + devices.register_device( + token="b" * 64, + bundle_id="org.solpbc.solstone-swift", + environment="development", + platform="ios", + ) + muted_facet = journal_copy / "facets" / "muted" + muted_facet.mkdir(parents=True, exist_ok=True) + (muted_facet / "facet.json").write_text( + json.dumps({"muted": True}), encoding="utf-8" + ) + _seed_activity( + journal_copy, + "montague", + "20260327", + [ + { + "id": "anticipated_meeting_090000_0327", + "source": "anticipated", + "start": "09:00", + "title": "Launch sync", + } + ], + ) + _seed_activity( + journal_copy, + "muted", + "20260327", + [ + { + "id": "anticipated_meeting_090000_muted", + "source": "anticipated", + "start": "09:00", + "title": "Muted meeting", + } + ], + ) + captured: list[str] = [] + + async def fake_post(self, url, *, headers, json): + captured.append(headers["apns-collapse-id"]) + return httpx.Response(200) + + with patch.object(httpx.AsyncClient, "post", new=fake_post): + triggers.check_pre_meeting_prep(datetime(2026, 3, 27, 8, 45, 0)) + + log_path = journal_copy / "push" / "nudge_log.jsonl" + lines = [ + json.loads(line) + for line in log_path.read_text(encoding="utf-8").splitlines() + if line.strip() + ] + assert captured == ["meeting.anticipated_meeting_090000_0327"] + assert any(line["category"] == "SOLSTONE_PRE_MEETING_PREP" for line in lines) diff --git a/tests/test_push_routes.py b/tests/test_push_routes.py new file mode 100644 index 000000000..5138d9c33 --- /dev/null +++ b/tests/test_push_routes.py @@ -0,0 +1,139 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import pytest + +from convey import create_app +from think.push.runtime import stop_all_push_runtime + + +@pytest.fixture +def push_app(journal_copy): + app = create_app(str(journal_copy)) + app.config["TESTING"] = True + yield app + stop_all_push_runtime() + + +@pytest.fixture +def push_client(push_app): + return push_app.test_client() + + +def test_register_push_device_happy_path(push_client, monkeypatch): + monkeypatch.setattr("convey.push.register_device", lambda **kwargs: 2) + + response = push_client.post( + "/api/push/register", + json={ + "device_token": "A" * 64, + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios", + }, + ) + + assert response.status_code == 200 + assert response.get_json() == {"registered": True, "device_count": 2} + + +def test_register_push_device_rejects_non_object(push_client): + response = push_client.post("/api/push/register", json=["bad"]) + + assert response.status_code == 400 + assert response.get_json() == {"error": "request body must be a JSON object"} + + +def test_register_push_device_validates_fields(push_client): + response = push_client.post("/api/push/register", json={"device_token": "x"}) + + assert response.status_code == 400 + assert response.get_json() == {"error": "bundle_id is required"} + + +def test_delete_push_device_happy_path(push_client, monkeypatch): + monkeypatch.setattr("convey.push.remove_device", lambda token: True) + monkeypatch.setattr("convey.push.load_devices", lambda: [{"token": "a"}]) + + response = push_client.delete("/api/push/register", json={"device_token": "a" * 64}) + + assert response.status_code == 200 + assert response.get_json() == {"removed": True, "device_count": 1} + + +def test_status_masks_tokens(push_client, monkeypatch): + monkeypatch.setattr("convey.push.is_configured", lambda: True) + monkeypatch.setattr( + "convey.push.load_devices", + lambda: [ + { + "token": "a" * 64, + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios", + "registered_at": 2, + } + ], + ) + monkeypatch.setattr( + "convey.push.status_view", + lambda device: { + "token_suffix": "...aaaa", + "bundle_id": device["bundle_id"], + "environment": device["environment"], + "platform": device["platform"], + "registered_at": "2024-04-19T12:00:00Z", + }, + ) + + response = push_client.get("/api/push/status") + + assert response.status_code == 200 + assert response.get_json() == { + "configured": True, + "device_count": 1, + "devices": [ + { + "token_suffix": "...aaaa", + "bundle_id": "org.solpbc.solstone-swift", + "environment": "development", + "platform": "ios", + "registered_at": "2024-04-19T12:00:00Z", + } + ], + } + + +def test_push_test_requires_configuration(push_client, monkeypatch): + monkeypatch.setattr("convey.push.is_configured", lambda: False) + + response = push_client.post("/api/push/test") + + assert response.status_code == 503 + assert response.get_json() == {"error": "push not configured"} + + +def test_push_test_validates_category(push_client, monkeypatch): + monkeypatch.setattr("convey.push.is_configured", lambda: True) + + response = push_client.post("/api/push/test", json={"category": "BAD"}) + + assert response.status_code == 400 + assert response.get_json() == {"error": "category must be a known push category"} + + +def test_push_test_happy_path(push_client, monkeypatch): + monkeypatch.setattr("convey.push.is_configured", lambda: True) + monkeypatch.setattr( + "convey.push.triggers.send_agent_alert", + lambda *, title, body, context_id: (1, 0), + ) + + response = push_client.post( + "/api/push/test", json={"title": "Alert", "body": "Body"} + ) + + assert response.status_code == 200 + assert response.get_json() == {"sent": 1, "failed": 0} diff --git a/tests/test_push_runtime.py b/tests/test_push_runtime.py new file mode 100644 index 000000000..1652bd3dd --- /dev/null +++ b/tests/test_push_runtime.py @@ -0,0 +1,95 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import pytest +from flask import Flask + +from think.push.runtime import ( + get_runtime_state, + start_push_runtime, + stop_all_push_runtime, + stop_push_runtime, +) + + +@pytest.fixture(autouse=True) +def reset_runtime(): + stop_all_push_runtime() + yield + stop_all_push_runtime() + + +def test_start_push_runtime_attaches_state(monkeypatch): + calls: list[str] = [] + monkeypatch.setattr( + "think.push.runtime.CallosumConnection.start", + lambda self, callback=None: calls.append("start"), + ) + monkeypatch.setattr( + "think.push.runtime.CallosumConnection.stop", + lambda self: calls.append("stop"), + ) + app = Flask(__name__) + + start_push_runtime(app) + try: + runtime = get_runtime_state() + assert app.push_runtime_started is True + assert runtime is not None + assert runtime.loop is not None + assert runtime.thread is not None + assert calls == ["start"] + finally: + stop_push_runtime(app) + + +def test_start_push_runtime_is_idempotent(monkeypatch): + monkeypatch.setattr( + "think.push.runtime.CallosumConnection.start", lambda self, callback=None: None + ) + monkeypatch.setattr("think.push.runtime.CallosumConnection.stop", lambda self: None) + app = Flask(__name__) + + start_push_runtime(app) + runtime = get_runtime_state() + first_loop = runtime.loop if runtime else None + first_thread = runtime.thread if runtime else None + try: + start_push_runtime(app) + runtime = get_runtime_state() + assert runtime is not None + assert runtime.loop is first_loop + assert runtime.thread is first_thread + assert runtime.apps.count(app) == 1 + finally: + stop_push_runtime(app) + + +def test_stop_push_runtime_cleans_last_app(monkeypatch): + monkeypatch.setattr( + "think.push.runtime.CallosumConnection.start", lambda self, callback=None: None + ) + monkeypatch.setattr("think.push.runtime.CallosumConnection.stop", lambda self: None) + app = Flask(__name__) + + start_push_runtime(app) + stop_push_runtime(app) + + assert app.push_runtime_started is False + assert get_runtime_state() is None + + +def test_stop_all_push_runtime_clears_runtime(monkeypatch): + monkeypatch.setattr( + "think.push.runtime.CallosumConnection.start", lambda self, callback=None: None + ) + monkeypatch.setattr("think.push.runtime.CallosumConnection.stop", lambda self: None) + app = Flask(__name__) + + start_push_runtime(app) + stop_all_push_runtime() + + assert app.push_runtime_started is False + assert get_runtime_state() is None diff --git a/tests/test_push_triggers.py b/tests/test_push_triggers.py new file mode 100644 index 000000000..aecc3c542 --- /dev/null +++ b/tests/test_push_triggers.py @@ -0,0 +1,212 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +from __future__ import annotations + +import json +from datetime import datetime +from pathlib import Path + +from think.push import triggers + + +def _log_path(tmp_path: Path) -> Path: + return tmp_path / "push" / "nudge_log.jsonl" + + +def test_handle_briefing_finish_polls_until_briefing_exists(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + responses = iter( + [ + ({}, None, []), + ({}, None, []), + ( + {"needs_attention": "- item"}, + {"generated": "2026-04-19T06:45:00"}, + ["one"], + ), + ] + ) + sent_calls: list[dict[str, object]] = [] + monkeypatch.setattr(triggers, "_load_briefing_md", lambda today: next(responses)) + monkeypatch.setattr(triggers.time, "sleep", lambda seconds: None) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: [{"token": "a" * 64}]) + monkeypatch.setattr( + triggers, + "send_many", + lambda devices, payload, *, collapse_id: ( + sent_calls.append( + {"devices": devices, "payload": payload, "collapse_id": collapse_id} + ) + or (1, 0) + ), + ) + + triggers.handle_briefing_finish( + {"tract": "cortex", "event": "finish", "name": "morning_briefing"} + ) + + assert len(sent_calls) == 1 + assert sent_calls[0]["collapse_id"].startswith("briefing.") + log_lines = _log_path(tmp_path).read_text(encoding="utf-8").splitlines() + assert len(log_lines) == 1 + + +def test_handle_briefing_finish_is_idempotent(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + sent_calls: list[str] = [] + monkeypatch.setattr( + triggers, + "_load_briefing_md", + lambda today: ( + {"needs_attention": "- item"}, + {"generated": "2026-04-19T06:45:00"}, + ["one"], + ), + ) + monkeypatch.setattr(triggers.time, "sleep", lambda seconds: None) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: [{"token": "a" * 64}]) + monkeypatch.setattr( + triggers, + "send_many", + lambda devices, payload, *, collapse_id: ( + sent_calls.append(collapse_id) or (1, 0) + ), + ) + + message = {"tract": "cortex", "event": "finish", "name": "morning_briefing"} + triggers.handle_briefing_finish(message) + triggers.handle_briefing_finish(message) + + assert sent_calls == [sent_calls[0]] + + +def test_check_pre_meeting_prep_skips_muted_facets(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: [{"token": "a" * 64}]) + monkeypatch.setattr(triggers, "get_enabled_facets", lambda: {}) + sent_calls: list[str] = [] + monkeypatch.setattr( + triggers, + "send_many", + lambda devices, payload, *, collapse_id: ( + sent_calls.append(collapse_id) or (1, 0) + ), + ) + + triggers.check_pre_meeting_prep(datetime(2026, 4, 20, 8, 45, 0)) + + assert sent_calls == [] + + +def test_check_pre_meeting_prep_skips_non_anticipated(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: [{"token": "a" * 64}]) + monkeypatch.setattr(triggers, "get_enabled_facets", lambda: {"work": {}}) + monkeypatch.setattr( + triggers, + "load_activity_records", + lambda facet, day: [{"id": "meeting", "source": "cogitate", "start": "09:00"}], + ) + sent_calls: list[str] = [] + monkeypatch.setattr( + triggers, + "send_many", + lambda devices, payload, *, collapse_id: ( + sent_calls.append(collapse_id) or (1, 0) + ), + ) + + triggers.check_pre_meeting_prep(datetime(2026, 4, 20, 8, 45, 0)) + + assert sent_calls == [] + + +def test_check_pre_meeting_prep_fires_for_hhmm_and_hhmmss(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: [{"token": "a" * 64}]) + monkeypatch.setattr(triggers, "get_enabled_facets", lambda: {"work": {}}) + monkeypatch.setattr( + triggers, + "load_activity_records", + lambda facet, day: [ + { + "id": "anticipated_meeting_090000_0420", + "source": "anticipated", + "start": "09:00", + "title": "Launch sync", + }, + { + "id": "anticipated_call_090030_0420", + "source": "anticipated", + "start": "09:00:30", + "title": "Prep call", + }, + ], + ) + sent_calls: list[str] = [] + monkeypatch.setattr( + triggers, + "send_many", + lambda devices, payload, *, collapse_id: ( + sent_calls.append(collapse_id) or (1, 0) + ), + ) + + triggers.check_pre_meeting_prep(datetime(2026, 4, 20, 8, 45, 0)) + + assert sent_calls == [ + "meeting.anticipated_meeting_090000_0420", + "meeting.anticipated_call_090030_0420", + ] + + +def test_check_pre_meeting_prep_zero_devices_skips_log(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: []) + monkeypatch.setattr(triggers, "get_enabled_facets", lambda: {"work": {}}) + monkeypatch.setattr( + triggers, + "load_activity_records", + lambda facet, day: [ + { + "id": "anticipated_meeting_090000_0420", + "source": "anticipated", + "start": "09:00", + "title": "Launch sync", + } + ], + ) + + triggers.check_pre_meeting_prep(datetime(2026, 4, 20, 8, 45, 0)) + + assert _log_path(tmp_path).exists() is False + + +def test_send_agent_alert_same_context_id_fires_once(monkeypatch, tmp_path): + monkeypatch.setenv("_SOLSTONE_JOURNAL_OVERRIDE", str(tmp_path)) + monkeypatch.setattr(triggers, "_eligible_devices", lambda: [{"token": "a" * 64}]) + sent_calls: list[str] = [] + monkeypatch.setattr( + triggers, + "send_many", + lambda devices, payload, *, collapse_id: ( + sent_calls.append(collapse_id) or (1, 0) + ), + ) + + first = triggers.send_agent_alert( + title="Agent Alert", body="Needs review", context_id="ctx-1" + ) + second = triggers.send_agent_alert( + title="Agent Alert", body="Needs review", context_id="ctx-1" + ) + + assert first == (1, 0) + assert second == (0, 0) + assert sent_calls == ["alert.ctx-1"] + lines = [ + json.loads(line) + for line in _log_path(tmp_path).read_text(encoding="utf-8").splitlines() + ] + assert len(lines) == 1 diff --git a/think/journal_default.json b/think/journal_default.json index d643df240..f00874d17 100644 --- a/think/journal_default.json +++ b/think/journal_default.json @@ -37,6 +37,13 @@ "model": "gpt-realtime", "brain_model": "haiku" }, + "push": { + "apns_key_path": null, + "apns_key_id": null, + "apns_team_id": null, + "bundle_id": null, + "environment": "development" + }, "retention": { "raw_media": "days", "raw_media_days": 7, diff --git a/think/push/__init__.py b/think/push/__init__.py new file mode 100644 index 000000000..774bb2b9f --- /dev/null +++ b/think/push/__init__.py @@ -0,0 +1,20 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Push package.""" + +from think.push.runtime import ( + get_runtime_state, + start_push_runtime, + stop_all_push_runtime, + stop_push_runtime, +) +from think.push.triggers import send_agent_alert + +__all__ = [ + "get_runtime_state", + "send_agent_alert", + "start_push_runtime", + "stop_all_push_runtime", + "stop_push_runtime", +] diff --git a/think/push/config.py b/think/push/config.py new file mode 100644 index 000000000..e9fa0f1d0 --- /dev/null +++ b/think/push/config.py @@ -0,0 +1,92 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Push config readers.""" + +from __future__ import annotations + +from pathlib import Path +from typing import Any + +from think.utils import get_config + +DEFAULT_ENVIRONMENT = "development" +_VALID_ENVIRONMENTS = {"development", "production"} + + +def _push_config() -> dict[str, Any]: + config = get_config() + push = config.get("push") + return push if isinstance(push, dict) else {} + + +def _clean_str(value: Any) -> str | None: + if not isinstance(value, str): + return None + cleaned = value.strip() + return cleaned or None + + +def get_apns_key_path() -> Path | None: + configured = _clean_str(_push_config().get("apns_key_path")) + return Path(configured) if configured else None + + +def get_apns_key_id() -> str | None: + return _clean_str(_push_config().get("apns_key_id")) + + +def get_apns_team_id() -> str | None: + return _clean_str(_push_config().get("apns_team_id")) + + +def get_bundle_id() -> str | None: + return _clean_str(_push_config().get("bundle_id")) + + +def get_environment() -> str: + configured = _clean_str(_push_config().get("environment")) + if configured is None: + return DEFAULT_ENVIRONMENT + if configured not in _VALID_ENVIRONMENTS: + raise ValueError( + "push.environment must be 'development' or 'production' when set" + ) + return configured + + +def _has_valid_key_path() -> bool: + key_path = get_apns_key_path() + if key_path is None or not key_path.is_absolute() or not key_path.is_file(): + return False + try: + key_path.read_text(encoding="utf-8") + except OSError: + return False + return True + + +def is_configured() -> bool: + if not ( + get_apns_key_path() + and get_apns_key_id() + and get_apns_team_id() + and get_bundle_id() + ): + return False + try: + get_environment() + except ValueError: + return False + return _has_valid_key_path() + + +__all__ = [ + "DEFAULT_ENVIRONMENT", + "get_apns_key_id", + "get_apns_key_path", + "get_apns_team_id", + "get_bundle_id", + "get_environment", + "is_configured", +] diff --git a/think/push/devices.py b/think/push/devices.py new file mode 100644 index 000000000..cb5871e85 --- /dev/null +++ b/think/push/devices.py @@ -0,0 +1,153 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Push device storage.""" + +from __future__ import annotations + +import json +import logging +import time +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +from think.entities.core import atomic_write +from think.utils import get_journal + +logger = logging.getLogger("solstone.push.devices") + + +def _devices_path() -> Path: + return Path(get_journal()) / "config" / "push_devices.json" + + +def _empty_store() -> dict[str, list[dict[str, Any]]]: + return {"devices": []} + + +def _validate_store(payload: Any) -> list[dict[str, Any]]: + if not isinstance(payload, dict): + raise ValueError("push device store must be a JSON object") + devices = payload.get("devices") + if not isinstance(devices, list): + raise ValueError("push device store must contain a devices list") + normalized: list[dict[str, Any]] = [] + for device in devices: + if not isinstance(device, dict): + raise ValueError("push device rows must be JSON objects") + token = str(device.get("token") or "").strip() + bundle_id = str(device.get("bundle_id") or "").strip() + environment = str(device.get("environment") or "").strip() + platform = str(device.get("platform") or "").strip() + registered_at = device.get("registered_at") + if ( + not token + or not bundle_id + or not environment + or not platform + or not isinstance(registered_at, (int, float)) + ): + raise ValueError("push device row missing required fields") + normalized.append( + { + "token": token, + "bundle_id": bundle_id, + "environment": environment, + "platform": platform, + "registered_at": int(registered_at), + } + ) + return normalized + + +def _read_store() -> list[dict[str, Any]]: + path = _devices_path() + if not path.exists(): + return [] + try: + payload = json.loads(path.read_text(encoding="utf-8")) + return _validate_store(payload) + except Exception as exc: + logger.warning("push device store unreadable path=%s error=%s", path, exc) + return [] + + +def _write_store(devices: list[dict[str, Any]]) -> None: + payload = json.dumps({"devices": devices}, indent=2, ensure_ascii=False) + "\n" + atomic_write(_devices_path(), payload, prefix=".push_devices_") + + +def load_devices() -> list[dict[str, Any]]: + return _read_store() + + +def register_device( + *, token: str, bundle_id: str, environment: str, platform: str +) -> int: + devices = load_devices() + registered_at = int(time.time()) + updated = False + for device in devices: + if device["token"] != token: + continue + device.update( + { + "bundle_id": bundle_id, + "environment": environment, + "platform": platform, + "registered_at": registered_at, + } + ) + updated = True + break + if not updated: + devices.append( + { + "token": token, + "bundle_id": bundle_id, + "environment": environment, + "platform": platform, + "registered_at": registered_at, + } + ) + _write_store(devices) + return len(devices) + + +def remove_device(token: str) -> bool: + devices = load_devices() + remaining = [device for device in devices if device["token"] != token] + if len(remaining) == len(devices): + return False + _write_store(remaining) + return True + + +def mask_token(token: str) -> str: + return "..." + str(token or "")[-4:] + + +def status_view(device: dict[str, Any]) -> dict[str, Any]: + registered_at = int(device["registered_at"]) + registered_at_label = ( + datetime.fromtimestamp(registered_at, tz=timezone.utc) + .isoformat() + .replace("+00:00", "Z") + ) + return { + "token_suffix": mask_token(device.get("token", "")), + "bundle_id": device["bundle_id"], + "environment": device["environment"], + "platform": device["platform"], + "registered_at": registered_at_label, + } + + +__all__ = [ + "load_devices", + "mask_token", + "register_device", + "remove_device", + "status_view", +] diff --git a/think/push/dispatch.py b/think/push/dispatch.py new file mode 100644 index 000000000..ad1290a9f --- /dev/null +++ b/think/push/dispatch.py @@ -0,0 +1,378 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""APNs transport for push notifications.""" + +from __future__ import annotations + +import asyncio +import json +import logging +import threading +import time +from pathlib import Path +from typing import Any + +import httpx +import jwt + +from think.push import devices +from think.push.config import ( + get_apns_key_id, + get_apns_key_path, + get_apns_team_id, + get_bundle_id, + get_environment, +) + +logger = logging.getLogger("solstone.push.dispatch") + +CATEGORY_DAILY_BRIEFING = "SOLSTONE_DAILY_BRIEFING" +CATEGORY_PRE_MEETING_PREP = "SOLSTONE_PRE_MEETING_PREP" +CATEGORY_AGENT_ALERT = "SOLSTONE_AGENT_ALERT" +CATEGORY_COMMITMENT_NUDGE = "SOLSTONE_COMMITMENT_NUDGE" +CATEGORIES = ( + CATEGORY_DAILY_BRIEFING, + CATEGORY_PRE_MEETING_PREP, + CATEGORY_AGENT_ALERT, + CATEGORY_COMMITMENT_NUDGE, +) +_JWT_MAX_AGE_SECONDS = 55 * 60 +_APNS_JWT_CACHE: dict[tuple[str, str], tuple[str, int]] = {} +_APNS_JWT_CACHE_LOCK = threading.Lock() + + +def _require_bundle_id() -> str: + bundle_id = get_bundle_id() + if bundle_id is None: + raise RuntimeError("push.bundle_id is not configured") + return bundle_id + + +def _require_key_id() -> str: + key_id = get_apns_key_id() + if key_id is None: + raise RuntimeError("push.apns_key_id is not configured") + return key_id + + +def _require_team_id() -> str: + team_id = get_apns_team_id() + if team_id is None: + raise RuntimeError("push.apns_team_id is not configured") + return team_id + + +def _require_key_path() -> Path: + key_path = get_apns_key_path() + if key_path is None: + raise RuntimeError("push.apns_key_path is not configured") + if not key_path.is_absolute(): + raise ValueError("push.apns_key_path must be an absolute path") + if not key_path.exists(): + raise FileNotFoundError(f"APNs key file not found: {key_path}") + if not key_path.is_file(): + raise RuntimeError(f"APNs key file is not a regular file: {key_path}") + return key_path + + +def _mint_apns_jwt(*, now: int | None = None) -> str: + issued_at = int(time.time()) if now is None else now + key_id = _require_key_id() + team_id = _require_team_id() + cache_key = (key_id, team_id) + with _APNS_JWT_CACHE_LOCK: + cached = _APNS_JWT_CACHE.get(cache_key) + if cached and issued_at - cached[1] <= _JWT_MAX_AGE_SECONDS: + return cached[0] + token = jwt.encode( + {"iss": team_id, "iat": issued_at}, + _require_key_path().read_text(encoding="utf-8"), + algorithm="ES256", + headers={"alg": "ES256", "kid": key_id}, + ) + _APNS_JWT_CACHE[cache_key] = (token, issued_at) + return token + + +def build_daily_briefing_collapse_id(day: str) -> str: + return f"briefing.{day}" + + +def build_pre_meeting_collapse_id(activity_id: str) -> str: + return f"meeting.{activity_id}" + + +def build_agent_alert_collapse_id(context_id: str) -> str: + return f"alert.{context_id}" + + +def build_commitment_collapse_id(ledger_id: str) -> str: + return f"commitment.{ledger_id}" + + +def build_daily_briefing_payload( + *, day: str, generated: str | None, needs_attention_count: int +) -> dict[str, Any]: + return { + "aps": { + "alert": { + "title": "Daily Briefing", + "body": "Your briefing is ready — tap to view", + }, + "category": CATEGORY_DAILY_BRIEFING, + "sound": "default", + "mutable-content": 1, + "content-available": 1, + }, + "data": { + "action": "open_briefing", + "day": day, + "generated": generated, + "needs_attention_count": needs_attention_count, + }, + } + + +def build_pre_meeting_payload( + *, activity: dict[str, Any], facet: str, day: str +) -> dict[str, Any]: + participants = [ + str(entry.get("name") or "").strip() + for entry in activity.get("participation", []) + if isinstance(entry, dict) + and entry.get("role") == "attendee" + and str(entry.get("name") or "").strip() + ] + return { + "aps": { + "alert": { + "title": "Pre-Meeting Prep", + "body": "Meeting in 15 minutes — tap to view", + }, + "category": CATEGORY_PRE_MEETING_PREP, + "sound": "default", + "mutable-content": 1, + "content-available": 1, + "interruption-level": "time-sensitive", + }, + "data": { + "action": "open_pre_meeting", + "activity_id": str(activity.get("id") or ""), + "facet": facet, + "day": day, + "start": str(activity.get("start") or ""), + "title": str(activity.get("title") or ""), + "location": str(activity.get("location") or ""), + "participants": participants, + "prep_notes": str(activity.get("prep_notes") or ""), + }, + } + + +def build_agent_alert_payload( + *, title: str, body: str, context_id: str +) -> dict[str, Any]: + return { + "aps": { + "alert": {"title": title, "body": body}, + "category": CATEGORY_AGENT_ALERT, + "sound": "default", + "mutable-content": 1, + "content-available": 1, + }, + "data": {"action": "open_alert", "context_id": context_id}, + } + + +def build_commitment_payload(*, ledger_id: str) -> dict[str, Any]: + return { + "aps": { + "alert": { + "title": "Commitment Nudge", + "body": "A commitment needs attention — tap to view", + }, + "category": CATEGORY_COMMITMENT_NUDGE, + "sound": "default", + "mutable-content": 1, + "content-available": 1, + }, + "data": {"action": "open_commitment", "ledger_id": ledger_id}, + } + + +def _apns_host() -> str: + environment = get_environment() + if environment == "production": + return "https://api.push.apple.com" + return "https://api.sandbox.push.apple.com" + + +def _headers(*, collapse_id: str, priority: int) -> dict[str, str]: + return { + "apns-topic": _require_bundle_id(), + "apns-collapse-id": collapse_id, + "apns-priority": str(priority), + "apns-push-type": "alert", + "authorization": f"bearer {_mint_apns_jwt()}", + } + + +def _response_reason(response: httpx.Response) -> str | None: + try: + payload = response.json() + except (json.JSONDecodeError, ValueError): + return None + if not isinstance(payload, dict): + return None + reason = payload.get("reason") + return str(reason) if isinstance(reason, str) and reason else None + + +def _run_async(coro: Any) -> Any: + try: + asyncio.get_running_loop() + except RuntimeError: + return asyncio.run(coro) + result: dict[str, Any] = {} + error: dict[str, BaseException] = {} + + def runner() -> None: + try: + result["value"] = asyncio.run(coro) + except BaseException as exc: + error["value"] = exc + + thread = threading.Thread(target=runner, name="push-dispatch", daemon=True) + thread.start() + thread.join() + if "value" in error: + raise error["value"] + return result.get("value") + + +async def _send_with_client( + client: httpx.AsyncClient, + device: dict[str, Any], + payload: dict[str, Any], + *, + collapse_id: str, + priority: int, +) -> tuple[bool, str | None]: + token = str(device.get("token") or "") + masked_token = devices.mask_token(token) + try: + response = await client.post( + f"{_apns_host()}/3/device/{token}", + headers=_headers(collapse_id=collapse_id, priority=priority), + json=payload, + ) + except Exception as exc: + logger.warning("push delivery failed token=%s error=%s", masked_token, exc) + return False, str(exc) + reason = _response_reason(response) + if response.status_code == 200: + return True, None + if response.status_code == 410 or reason in {"BadDeviceToken", "Unregistered"}: + devices.remove_device(token) + logger.warning( + "push pruning token=%s status=%s reason=%s", + masked_token, + response.status_code, + reason or "", + ) + return False, reason + if 500 <= response.status_code: + logger.warning( + "push rejected token=%s status=%s reason=%s", + masked_token, + response.status_code, + reason or "", + ) + return False, reason + logger.error( + "push rejected token=%s status=%s reason=%s", + masked_token, + response.status_code, + reason or "", + ) + return False, reason + + +async def _send_async( + device: dict[str, Any], + payload: dict[str, Any], + *, + collapse_id: str, + priority: int = 10, +) -> tuple[bool, str | None]: + async with httpx.AsyncClient(http2=True, timeout=10.0) as client: + return await _send_with_client( + client, device, payload, collapse_id=collapse_id, priority=priority + ) + + +def send( + device: dict[str, Any], + payload: dict[str, Any], + *, + collapse_id: str, + priority: int = 10, +) -> tuple[bool, str | None]: + return _run_async( + _send_async(device, payload, collapse_id=collapse_id, priority=priority) + ) + + +async def _send_many_async( + push_devices: list[dict[str, Any]], + payload: dict[str, Any], + *, + collapse_id: str, + priority: int = 10, +) -> tuple[int, int]: + sent = 0 + failed = 0 + async with httpx.AsyncClient(http2=True, timeout=10.0) as client: + for device in push_devices: + ok, _ = await _send_with_client( + client, device, payload, collapse_id=collapse_id, priority=priority + ) + if ok: + sent += 1 + else: + failed += 1 + return sent, failed + + +def send_many( + push_devices: list[dict[str, Any]], + payload: dict[str, Any], + *, + collapse_id: str, + priority: int = 10, +) -> tuple[int, int]: + return _run_async( + _send_many_async( + push_devices, payload, collapse_id=collapse_id, priority=priority + ) + ) + + +__all__ = [ + "CATEGORIES", + "CATEGORY_AGENT_ALERT", + "CATEGORY_COMMITMENT_NUDGE", + "CATEGORY_DAILY_BRIEFING", + "CATEGORY_PRE_MEETING_PREP", + "build_agent_alert_collapse_id", + "build_agent_alert_payload", + "build_commitment_collapse_id", + "build_commitment_payload", + "build_daily_briefing_collapse_id", + "build_daily_briefing_payload", + "build_pre_meeting_collapse_id", + "build_pre_meeting_payload", + "send", + "send_many", +] diff --git a/think/push/runtime.py b/think/push/runtime.py new file mode 100644 index 000000000..3b4fba466 --- /dev/null +++ b/think/push/runtime.py @@ -0,0 +1,144 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Background runtime for push tasks.""" + +from __future__ import annotations + +import asyncio +import atexit +import logging +import threading +from dataclasses import dataclass, field +from datetime import datetime +from typing import Any + +from think.callosum import CallosumConnection +from think.push import triggers + +logger = logging.getLogger("solstone.push.runtime") + + +@dataclass +class RuntimeState: + loop: asyncio.AbstractEventLoop | None = None + thread: threading.Thread | None = None + started_event: threading.Event = field(default_factory=threading.Event) + apps: list[Any] = field(default_factory=list) + callosum: CallosumConnection | None = None + periodic_task: asyncio.Task[Any] | None = None + + +_RUNTIME_LOCK = threading.Lock() +_runtime: RuntimeState | None = None +_atexit_registered = False + + +def get_runtime_state() -> RuntimeState | None: + return _runtime + + +def _on_callosum_message(message: dict[str, Any]) -> None: + try: + triggers.handle_briefing_finish(message) + except Exception: + logger.exception("push callosum handler failed") + + +async def _periodic_loop() -> None: + while True: + await asyncio.sleep(60) + try: + triggers.check_pre_meeting_prep(datetime.now()) + except Exception: + logger.exception("push periodic check failed") + + +def _thread_main(runtime: RuntimeState) -> None: + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + runtime.loop = loop + runtime.callosum = CallosumConnection() + runtime.callosum.start(callback=_on_callosum_message) + runtime.periodic_task = loop.create_task(_periodic_loop()) + runtime.started_event.set() + try: + loop.run_forever() + finally: + pending = [task for task in asyncio.all_tasks(loop) if not task.done()] + for task in pending: + task.cancel() + if pending: + loop.run_until_complete(asyncio.gather(*pending, return_exceptions=True)) + loop.close() + + +def start_push_runtime(app: Any) -> None: + global _runtime, _atexit_registered + + with _RUNTIME_LOCK: + if _runtime is None: + runtime = RuntimeState() + thread = threading.Thread( + target=_thread_main, + args=(runtime,), + name="push-runtime", + daemon=True, + ) + runtime.thread = thread + _runtime = runtime + thread.start() + runtime = _runtime + if app not in runtime.apps: + runtime.apps.append(app) + app.push_runtime_started = True + if not _atexit_registered: + atexit.register(stop_all_push_runtime) + _atexit_registered = True + started_event = runtime.started_event + started_event.wait(timeout=1.0) + + +def stop_push_runtime(app: Any) -> None: + runtime = _runtime + app.push_runtime_started = False + if runtime is None: + return + with _RUNTIME_LOCK: + if app in runtime.apps: + runtime.apps.remove(app) + remaining = list(runtime.apps) + if not remaining: + stop_all_push_runtime() + + +def stop_all_push_runtime() -> None: + global _runtime + + with _RUNTIME_LOCK: + runtime = _runtime + _runtime = None + if runtime is None: + return + for app in list(runtime.apps): + try: + app.push_runtime_started = False + except Exception: + logger.exception("push runtime app cleanup failed") + if runtime.callosum is not None: + runtime.callosum.stop() + if runtime.loop is not None: + if runtime.periodic_task is not None: + runtime.loop.call_soon_threadsafe(runtime.periodic_task.cancel) + runtime.loop.call_soon_threadsafe(runtime.loop.stop) + if runtime.thread is not None: + runtime.thread.join(timeout=1.0) + + +__all__ = [ + "RuntimeState", + "get_runtime_state", + "start_push_runtime", + "stop_all_push_runtime", + "stop_push_runtime", +] diff --git a/think/push/triggers.py b/think/push/triggers.py new file mode 100644 index 000000000..917a454c9 --- /dev/null +++ b/think/push/triggers.py @@ -0,0 +1,240 @@ +# SPDX-License-Identifier: AGPL-3.0-only +# Copyright (c) 2026 sol pbc + +"""Push trigger handlers.""" + +from __future__ import annotations + +import json +import logging +import time +from datetime import datetime +from pathlib import Path +from typing import Any + +from apps.home.routes import _load_briefing_md +from think.activities import load_activity_records +from think.facets import get_enabled_facets +from think.push.config import get_bundle_id, get_environment, is_configured +from think.push.devices import load_devices +from think.push.dispatch import ( + CATEGORY_AGENT_ALERT, + CATEGORY_DAILY_BRIEFING, + CATEGORY_PRE_MEETING_PREP, + build_agent_alert_collapse_id, + build_agent_alert_payload, + build_daily_briefing_collapse_id, + build_daily_briefing_payload, + build_pre_meeting_collapse_id, + build_pre_meeting_payload, + send_many, +) +from think.utils import get_journal + +logger = logging.getLogger("solstone.push.triggers") + + +def _nudge_log_path() -> Path: + return Path(get_journal()) / "push" / "nudge_log.jsonl" + + +def _serialize_dedupe_key(dedupe_key: tuple[Any, ...]) -> str: + return json.dumps(list(dedupe_key), separators=(",", ":"), ensure_ascii=False) + + +def _has_nudged(dedupe_key: tuple[Any, ...]) -> bool: + path = _nudge_log_path() + if not path.exists(): + return False + encoded = _serialize_dedupe_key(dedupe_key) + for line in path.read_text(encoding="utf-8").splitlines(): + if not line.strip(): + continue + try: + payload = json.loads(line) + except json.JSONDecodeError: + continue + if isinstance(payload, dict) and payload.get("dedupe_key") == encoded: + return True + return False + + +def _append_nudge_log(line: dict[str, Any]) -> None: + path = _nudge_log_path() + path.parent.mkdir(parents=True, exist_ok=True) + with path.open("a", encoding="utf-8") as handle: + handle.write(json.dumps(line, ensure_ascii=False) + "\n") + + +def _eligible_devices() -> list[dict[str, Any]]: + if not is_configured(): + logger.debug("push skipped configured=false") + return [] + bundle_id = get_bundle_id() + environment = get_environment() + matched = [ + device + for device in load_devices() + if device.get("bundle_id") == bundle_id + and device.get("environment") == environment + and device.get("platform") == "ios" + ] + if not matched: + logger.debug("push skipped devices=0") + return matched + + +def _metadata_generated(metadata: dict[str, Any] | None) -> str | None: + if not isinstance(metadata, dict): + return None + generated = metadata.get("generated") + if isinstance(generated, str): + return generated + if hasattr(generated, "isoformat"): + return generated.isoformat() + return None + + +def _record_send( + *, + dedupe_key: tuple[Any, ...], + category: str, + sent: int, + failed: int, + **payload: Any, +) -> None: + _append_nudge_log( + { + "ts": int(time.time()), + "category": category, + "dedupe_key": _serialize_dedupe_key(dedupe_key), + "sent": sent, + "failed": failed, + **payload, + } + ) + + +def handle_briefing_finish(message: dict[str, Any]) -> None: + if message.get("tract") != "cortex": + return + if message.get("event") != "finish": + return + if message.get("name") != "morning_briefing": + return + today = datetime.now().strftime("%Y%m%d") + dedupe_key = (CATEGORY_DAILY_BRIEFING, today) + if _has_nudged(dedupe_key): + return + sections: dict[str, str] = {} + metadata: dict[str, Any] | None = None + needs_attention: list[str] = [] + for _ in range(10): + sections, metadata, needs_attention = _load_briefing_md(today) + if sections and metadata: + break + time.sleep(1) + else: + logger.warning("push briefing unavailable after finish day=%s", today) + return + eligible_devices = _eligible_devices() + if not eligible_devices: + return + sent, failed = send_many( + eligible_devices, + build_daily_briefing_payload( + day=today, + generated=_metadata_generated(metadata), + needs_attention_count=len(needs_attention), + ), + collapse_id=build_daily_briefing_collapse_id(today), + ) + if sent > 0: + _record_send( + dedupe_key=dedupe_key, + category=CATEGORY_DAILY_BRIEFING, + day=today, + sent=sent, + failed=failed, + ) + + +def _parse_start(now: datetime, start: str) -> datetime | None: + for pattern in ("%H:%M", "%H:%M:%S"): + try: + parsed = datetime.strptime(start, pattern) + except ValueError: + continue + return now.replace( + hour=parsed.hour, + minute=parsed.minute, + second=parsed.second, + microsecond=0, + ) + return None + + +def check_pre_meeting_prep(now: datetime) -> None: + today = now.strftime("%Y%m%d") + eligible_devices = _eligible_devices() + if not eligible_devices: + return + for facet in get_enabled_facets().keys(): + for record in load_activity_records(facet, today): + if record.get("source") != "anticipated": + continue + activity_id = str(record.get("id") or "").strip() + start = str(record.get("start") or "").strip() + if not activity_id or not start: + continue + event_start = _parse_start(now, start) + if event_start is None: + logger.debug("push skipped invalid meeting start id=%s", activity_id) + continue + seconds_until = (event_start - now).total_seconds() + if seconds_until < 14 * 60 or seconds_until > 16 * 60: + continue + dedupe_key = (CATEGORY_PRE_MEETING_PREP, activity_id, today) + if _has_nudged(dedupe_key): + continue + sent, failed = send_many( + eligible_devices, + build_pre_meeting_payload(activity=record, facet=facet, day=today), + collapse_id=build_pre_meeting_collapse_id(activity_id), + ) + if sent > 0: + _record_send( + dedupe_key=dedupe_key, + category=CATEGORY_PRE_MEETING_PREP, + day=today, + facet=facet, + activity_id=activity_id, + sent=sent, + failed=failed, + ) + + +def send_agent_alert(*, title: str, body: str, context_id: str) -> tuple[int, int]: + dedupe_key = (CATEGORY_AGENT_ALERT, context_id) + if _has_nudged(dedupe_key): + return 0, 0 + eligible_devices = _eligible_devices() + if not eligible_devices: + return 0, 0 + sent, failed = send_many( + eligible_devices, + build_agent_alert_payload(title=title, body=body, context_id=context_id), + collapse_id=build_agent_alert_collapse_id(context_id), + ) + if sent > 0: + _record_send( + dedupe_key=dedupe_key, + category=CATEGORY_AGENT_ALERT, + context_id=context_id, + sent=sent, + failed=failed, + ) + return sent, failed + + +__all__ = ["check_pre_meeting_prep", "handle_briefing_finish", "send_agent_alert"] diff --git a/uv.lock b/uv.lock index ee3fad241..431b7ffb0 100644 --- a/uv.lock +++ b/uv.lock @@ -1009,6 +1009,19 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/04/4b/29cac41a4d98d144bf5f6d33995617b185d14b22401f75ca86f384e87ff1/h11-0.16.0-py3-none-any.whl", hash = "sha256:63cf8bbe7522de3bf65932fda1d9c2772064ffb3dae62d55932da54b31cb6c86", size = 37515, upload-time = "2025-04-24T03:35:24.344Z" }, ] +[[package]] +name = "h2" +version = "4.3.0" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "hpack" }, + { name = "hyperframe" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/1d/17/afa56379f94ad0fe8defd37d6eb3f89a25404ffc71d4d848893d270325fc/h2-4.3.0.tar.gz", hash = "sha256:6c59efe4323fa18b47a632221a1888bd7fde6249819beda254aeca909f221bf1", size = 2152026, upload-time = "2025-08-23T18:12:19.778Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/69/b2/119f6e6dcbd96f9069ce9a2665e0146588dc9f88f29549711853645e736a/h2-4.3.0-py3-none-any.whl", hash = "sha256:c438f029a25f7945c69e0ccf0fb951dc3f73a5f6412981daee861431b70e2bdd", size = 61779, upload-time = "2025-08-23T18:12:17.779Z" }, +] + [[package]] name = "hf-xet" version = "1.2.0" @@ -1038,6 +1051,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/cb/44/870d44b30e1dcfb6a65932e3e1506c103a8a5aea9103c337e7a53180322c/hf_xet-1.2.0-cp37-abi3-win_amd64.whl", hash = "sha256:e6584a52253f72c9f52f9e549d5895ca7a471608495c4ecaa6cc73dba2b24d69", size = 2905735, upload-time = "2025-10-24T19:04:35.928Z" }, ] +[[package]] +name = "hpack" +version = "4.1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/2c/48/71de9ed269fdae9c8057e5a4c0aa7402e8bb16f2c6e90b3aa53327b113f8/hpack-4.1.0.tar.gz", hash = "sha256:ec5eca154f7056aa06f196a557655c5b009b382873ac8d1e66e79e87535f1dca", size = 51276, upload-time = "2025-01-22T21:44:58.347Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/07/c6/80c95b1b2b94682a72cbdbfb85b81ae2daffa4291fbfa1b1464502ede10d/hpack-4.1.0-py3-none-any.whl", hash = "sha256:157ac792668d995c657d93111f46b4535ed114f0c9c8d672271bbec7eae1b496", size = 34357, upload-time = "2025-01-22T21:44:56.92Z" }, +] + [[package]] name = "httpcore" version = "1.0.9" @@ -1096,6 +1118,15 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/d5/ae/2f6d96b4e6c5478d87d606a1934b5d436c4a2bce6bb7c6fdece891c128e3/huggingface_hub-1.4.1-py3-none-any.whl", hash = "sha256:9931d075fb7a79af5abc487106414ec5fba2c0ae86104c0c62fd6cae38873d18", size = 553326, upload-time = "2026-02-06T09:20:00.728Z" }, ] +[[package]] +name = "hyperframe" +version = "6.1.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/02/e7/94f8232d4a74cc99514c13a9f995811485a6903d48e5d952771ef6322e30/hyperframe-6.1.0.tar.gz", hash = "sha256:f630908a00854a7adeabd6382b43923a4c4cd4b821fcb527e6ab9e15382a3b08", size = 26566, upload-time = "2025-01-22T21:41:49.302Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/48/30/47d0bf6072f7252e6521f3447ccfa40b421b6824517f82854703d0f5a98b/hyperframe-6.1.0-py3-none-any.whl", hash = "sha256:b03380493a519fce58ea5af42e4a42317bf9bd425596f7a0835ffce80f1a42e5", size = 13007, upload-time = "2025-01-22T21:41:47.295Z" }, +] + [[package]] name = "icalendar" version = "7.0.3" @@ -3538,6 +3569,7 @@ dependencies = [ { name = "freezegun" }, { name = "genai-prices" }, { name = "google-genai" }, + { name = "h2" }, { name = "httpx" }, { name = "icalendar" }, { name = "jsonschema" }, @@ -3553,6 +3585,7 @@ dependencies = [ { name = "pillow" }, { name = "playwright" }, { name = "psutil" }, + { name = "pyjwt" }, { name = "pyopenssl" }, { name = "pypdf" }, { name = "pytesseract" }, @@ -3588,6 +3621,7 @@ requires-dist = [ { name = "freezegun" }, { name = "genai-prices" }, { name = "google-genai" }, + { name = "h2" }, { name = "httpx" }, { name = "icalendar" }, { name = "jsonschema", specifier = ">=4.26,<5" }, @@ -3602,6 +3636,7 @@ requires-dist = [ { name = "pillow" }, { name = "playwright", specifier = ">=1.40.0" }, { name = "psutil" }, + { name = "pyjwt", specifier = ">=2.8" }, { name = "pyopenssl", specifier = ">=24.0" }, { name = "pypdf" }, { name = "pytesseract" },