Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 6 additions & 5 deletions database/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 8 additions & 0 deletions scripts/run_all.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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}"
Expand Down
10 changes: 3 additions & 7 deletions services/query_profiler_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
23 changes: 21 additions & 2 deletions webapp/app.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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,
Expand Down
115 changes: 89 additions & 26 deletions webapp/push_api.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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

Expand Down Expand Up @@ -493,31 +509,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.
Expand All @@ -532,14 +535,74 @@ 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)
Comment thread
amirbiron marked this conversation as resolved.
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
Expand Down
Loading