From b40a625fa6ffc7836d3b27f398bd3a04669c41f4 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 8 Jan 2026 11:34:53 +0000 Subject: [PATCH 1/4] feat: Cache public stats and optimize reminder polling Adds caching for the public stats API endpoint. Optimizes the note_reminders index and polling query for better performance. Co-authored-by: amirbiron --- database/manager.py | 11 ++++--- webapp/app.py | 23 ++++++++++++-- webapp/push_api.py | 77 ++++++++++++++++++++++++++++++--------------- 3 files changed, 78 insertions(+), 33 deletions(-) diff --git a/database/manager.py b/database/manager.py index 3899da86d..cce9925bd 100644 --- a/database/manager.py +++ b/database/manager.py @@ -772,16 +772,17 @@ def safe_create_index(*args: Any, **kwargs: Any) -> None: return DatabaseManager.safe_create_index(self, *args, **kwargs) # תיקון השגיאה ב-users: לא מבצעים בדיקה בוליאנית על Collection (PyMongo זורק חריגה) - # note_reminders - אינדקס מותאם לשאילתת הפולינג החדשה (push_api.py) - # השאילתה מסננת לפי: ack_at=null, status in [pending, snoozed], remind_at <= now, needs_push != false - # אינדקס חלקי (Partial Index) על ack_at=null ו-needs_push != false + # note_reminders - אינדקס מותאם לשאילתת הפולינג (push_api.py) + # קריטי: אינדקס חלקי לא תומך ב-$ne/$not ב-partialFilterExpression. + # לכן אנחנו מאנדקסים רק מסמכים "חדשים" עם needs_push=True. safe_create_index( "note_reminders", - [("status", ASCENDING), ("remind_at", ASCENDING), ("needs_push", ASCENDING)], + # remind_at ראשון כדי לאפשר sort+limit יעילים על due reminders + [("remind_at", ASCENDING), ("status", ASCENDING)], name="push_polling_optimized_idx", background=True, enforce=True, - partial_filter_expression={"ack_at": None, "needs_push": {"$ne": False}}, + partial_filter_expression={"ack_at": None, "needs_push": True}, ) # note_reminders - אינדקס לשאילתת /reminders/summary (sticky_notes_api.py) # השאילתה מסננת לפי: user_id, ack_at=null, status, remind_at diff --git a/webapp/app.py b/webapp/app.py index acf1656da..07ae8ebe3 100644 --- a/webapp/app.py +++ b/webapp/app.py @@ -14841,6 +14841,18 @@ def api_public_stats(): - total_snippets: סה"כ קטעי קוד ייחודיים שנשמרו אי פעם (distinct לפי user_id+file_name) כאשר התוכן לא ריק — כולל כאלה שנמחקו (is_active=false) """ try: + # קאש ל-5–10 דקות (באמצעות cache_manager; Redis אם זמין, אחרת פולבק בזיכרון). + # זה endpoint ציבורי שמרונדר הרבה ואין סיבה להפעיל aggregation כבד בכל ריענון. + cache_key = "api:public_stats:v1" + try: + cached = cache.get(cache_key) + if isinstance(cached, dict) and cached.get("ok") is True: + payload = dict(cached) + payload["cached"] = True + return jsonify(payload) + except Exception: + pass + db = get_db() now_utc = datetime.now(timezone.utc) last_24h = now_utc - timedelta(hours=24) @@ -14877,13 +14889,20 @@ def api_public_stats(): except Exception: total_snippets = 0 - return jsonify({ + payload = { "ok": True, "total_users": total_users, "active_users_24h": active_users_24h, "total_snippets": total_snippets, "timestamp": now_utc.isoformat(), - }) + "cached": False, + } + try: + # Dynamic TTL: public_stats ברירת מחדל 10 דקות (עם התאמות פעילות) + cache.set_dynamic(cache_key, payload, "public_stats", {"endpoint": "api_public_stats"}) + except Exception: + pass + return jsonify(payload) except Exception as e: return jsonify({ "ok": False, diff --git a/webapp/push_api.py b/webapp/push_api.py index 7f4ece83f..ac7eb522f 100644 --- a/webapp/push_api.py +++ b/webapp/push_api.py @@ -493,31 +493,18 @@ def _send_due_once(max_users: int = 100, max_per_user: int = 10) -> None: """ db = get_db() now = datetime.now(timezone.utc) - # Optimized query with backward compatibility for documents without needs_push. - # NOTE: our reminders use status pending/snoozed (not "active"). - mongo_filter = { + # IMPORTANT: do NOT use a single $or query here. + # A partial index (ack_at=None, needs_push=True) is only usable when the query + # *guarantees* needs_push=True. With a top-level $or (legacy branch), MongoDB + # tends to fall back to COLLSCAN under load. + # + # Strategy: + # 1) Query "new" documents (needs_push=True) — hits the partial index. + # 2) If we still need items, query legacy documents without needs_push. + base_filter = { "ack_at": None, "status": {"$in": ["pending", "snoozed"]}, "remind_at": {"$lte": now}, - "$or": [ - # New documents: explicit needs_push flag - {"needs_push": True}, - # Old documents without needs_push: use timestamp-based logic - # (never pushed, or remind_at changed since last push) - { - "needs_push": {"$exists": False}, - "$or": [ - {"last_push_success_at": {"$exists": False}}, - {"last_push_success_at": None}, - { - "$and": [ - {"last_push_success_at": {"$type": "date"}}, - {"$expr": {"$gt": ["$remind_at", "$last_push_success_at"]}}, - ] - }, - ], - }, - ], } total_needed = max_users * max_per_user # We still oversample to compensate for claim collisions / missing subscriptions, etc. @@ -532,14 +519,52 @@ def _send_due_once(max_users: int = 100, max_per_user: int = 10) -> None: "last_push_success_at": 1, } + raw_due: list = [] + # 1) New documents (fast path, uses partial index) try: - raw_due = list( - db.note_reminders.find(mongo_filter, projection).sort("remind_at", 1).limit(raw_limit) - ) + new_filter = dict(base_filter) + new_filter["needs_push"] = True + raw_due = list(db.note_reminders.find(new_filter, projection).sort("remind_at", 1).limit(raw_limit)) except Exception: raw_due = [] - due = [r for r in raw_due if isinstance(r, dict)] + due: list[dict] = [r for r in raw_due if isinstance(r, dict)] + + # 2) Legacy documents (best-effort, limited) + if len(due) < raw_limit: + try: + legacy_limit = max(1, raw_limit - len(due)) + legacy_filter = dict(base_filter) + legacy_filter["needs_push"] = {"$exists": False} + legacy_filter["$or"] = [ + {"last_push_success_at": {"$exists": False}}, + {"last_push_success_at": None}, + { + "$and": [ + {"last_push_success_at": {"$type": "date"}}, + {"$expr": {"$gt": ["$remind_at", "$last_push_success_at"]}}, + ] + }, + ] + legacy_raw = list( + db.note_reminders.find(legacy_filter, projection).sort("remind_at", 1).limit(legacy_limit) + ) + # Merge & de-dup by _id + seen = {str(d.get("_id")) for d in due if d.get("_id") is not None} + for r in legacy_raw: + if not isinstance(r, dict): + continue + rid = r.get("_id") + if rid is None: + continue + s = str(rid) + if s in seen: + continue + seen.add(s) + due.append(r) + except Exception: + pass + due = due[:total_needed] if not due: return From 9ad50f4bde819198a2d3d2d854b6d5825fbfbb68 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 8 Jan 2026 11:41:13 +0000 Subject: [PATCH 2/4] Fix: Ensure correct reminder order by sorting on remind_at Co-authored-by: amirbiron --- webapp/push_api.py | 22 ++++++++++++++++++++++ 1 file changed, 22 insertions(+) diff --git a/webapp/push_api.py b/webapp/push_api.py index ac7eb522f..5e5fa7507 100644 --- a/webapp/push_api.py +++ b/webapp/push_api.py @@ -565,6 +565,28 @@ def _send_due_once(max_users: int = 100, max_per_user: int = 10) -> None: except Exception: pass + # Keep global order by remind_at after merging new+legacy lists. + # Otherwise, early legacy reminders can be pushed after later "new" reminders. + def _due_sort_key(d: dict): + ra = d.get("remind_at") + if isinstance(ra, datetime): + try: + if ra.tzinfo is None: + ra = ra.replace(tzinfo=timezone.utc) + except Exception: + pass + else: + try: + ra = datetime.max.replace(tzinfo=timezone.utc) + except Exception: + ra = datetime.max + return (ra, str(d.get("_id") or "")) + + try: + due.sort(key=_due_sort_key) + except Exception: + pass + due = due[:total_needed] if not due: return From b73337d3e0b198a49b352e203df8d9f0b30577a6 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 8 Jan 2026 11:48:04 +0000 Subject: [PATCH 3/4] Refactor: Prevent duplicate push notification sender processes Co-authored-by: amirbiron --- scripts/run_all.sh | 8 ++++++++ webapp/push_api.py | 16 ++++++++++++++++ 2 files changed, 24 insertions(+) diff --git a/scripts/run_all.sh b/scripts/run_all.sh index 13ca29242..68a1321e8 100644 --- a/scripts/run_all.sh +++ b/scripts/run_all.sh @@ -14,6 +14,14 @@ is_true() { WEBAPP_START_SCRIPT="${WEBAPP_START_SCRIPT:-scripts/start_webapp.sh}" +# When running WebApp + AI service in the same container, keep Gunicorn concurrency conservative +# unless the operator explicitly overrides it. This reduces duplicate startup side-effects +# (jobs registration, background threads) and prevents double polling pressure on MongoDB. +if [[ -z "${WEB_CONCURRENCY:-}" && -z "${WEBAPP_GUNICORN_WORKERS:-}" ]]; then + export WEB_CONCURRENCY="1" + log "WEB_CONCURRENCY not set; defaulting to ${WEB_CONCURRENCY} (unified container)" +fi + # Internal AI Explain service (aiohttp) settings OBS_AI_EXPLAIN_RUN_LOCAL_SERVICE="${OBS_AI_EXPLAIN_RUN_LOCAL_SERVICE:-true}" OBS_AI_EXPLAIN_INTERNAL_HOST="${OBS_AI_EXPLAIN_INTERNAL_HOST:-127.0.0.1}" diff --git a/webapp/push_api.py b/webapp/push_api.py index 5e5fa7507..8da891c82 100644 --- a/webapp/push_api.py +++ b/webapp/push_api.py @@ -440,6 +440,7 @@ def unsubscribe(): # --- Background sender (opt-in via env) --- _sender_started = False +_sender_lock_fh = None # type: ignore def start_sender_if_enabled() -> None: @@ -449,6 +450,21 @@ def start_sender_if_enabled() -> None: enabled = (os.getenv("PUSH_NOTIFICATIONS_ENABLED", "true").strip().lower() in {"1", "true", "yes", "on"}) if not enabled: return + # Ensure only one sender loop runs across multiple Gunicorn workers/processes. + # Using an OS-level flock means the lock is released automatically on process exit/crash. + global _sender_lock_fh + try: + import fcntl + lock_path = os.getenv("PUSH_SENDER_LOCK_FILE") or "/tmp/codebot-push-sender.lock" + _sender_lock_fh = open(lock_path, "a+") + try: + fcntl.flock(_sender_lock_fh.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + except Exception: + # Another process owns the sender lock; do not start a duplicate loop. + return + except Exception: + # Fail-open: if locking isn't available, fall back to per-process behavior. + pass try: import threading From 2767c9375ef16e3010340c70ba93e6a99c244fd2 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Thu, 8 Jan 2026 11:55:26 +0000 Subject: [PATCH 4/4] Refactor query profiler to avoid duplicate logging Co-authored-by: amirbiron --- services/query_profiler_service.py | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/services/query_profiler_service.py b/services/query_profiler_service.py index 82cb184b4..90dad0bbf 100644 --- a/services/query_profiler_service.py +++ b/services/query_profiler_service.py @@ -429,13 +429,9 @@ def record_slow_query_sync( except Exception: pass - logger.warning( - "Slow query recorded: %s.%s took %.2fms (threshold: %sms)", - collection, - operation, - float(execution_time_ms), - self.slow_threshold_ms, - ) + # NOTE: We intentionally avoid a second free-form log line here. + # The structured JSON event above (emit_event) is the source of truth + # and prevents duplicate log lines in container logs. return record