Coverage for src/local_deep_research/news/flask_api.py: 93%

694 statements  

« prev     ^ index     » next       coverage.py v7.15.1, created at 2026-07-19 23:35 +0000

1""" 

2Flask API endpoints for news system. 

3Converted from FastAPI to match LDR's Flask architecture. 

4""" 

5 

6from functools import wraps 

7from typing import Any 

8from flask import Blueprint, request, jsonify 

9from loguru import logger 

10 

11from . import api 

12from .folder_manager import FolderManager 

13from ..database.models import SubscriptionFolder 

14from ..web.auth.decorators import login_required 

15from ..database.session_context import get_user_db_session 

16from ..settings.env_registry import get_env_setting 

17from ..utilities.db_utils import get_settings_manager 

18from ..llm.providers.base import normalize_provider 

19from ..security.decorators import require_json_body 

20 

21# Hard ceiling for user-supplied ``limit`` query params on the news 

22# endpoints in this module. Matches the ``max_value`` of the 

23# ``news.feed.default_limit`` setting so a direct API caller cannot request a 

24# larger page than the UI can configure. Shared by the feed and 

25# subscription-history endpoints here so their caps stay in lockstep. 

26NEWS_FEED_MAX_LIMIT = 100 

27 

28 

29def scheduler_control_required(f): 

30 """Decorator that gates global scheduler control behind a setting. 

31 

32 The news scheduler is a global singleton — starting, stopping, or 

33 triggering it affects all users. This decorator checks the 

34 ``news.scheduler.allow_api_control`` setting (env var 

35 ``LDR_NEWS_SCHEDULER_ALLOW_API_CONTROL``, default ``false``) and 

36 returns 403 when the setting is disabled. 

37 

38 Must be placed *after* ``@login_required`` in the decorator stack. 

39 """ 

40 

41 @wraps(f) 

42 def wrapper(*args, **kwargs): 

43 if not get_env_setting("news.scheduler.allow_api_control", False): 

44 from flask import session as flask_session 

45 

46 username = flask_session.get("username", "unknown") 

47 remote_addr = request.remote_addr 

48 logger.warning( 

49 "Scheduler API control blocked for endpoint {} (user={}, ip={})", 

50 f.__name__, 

51 username, 

52 remote_addr, 

53 ) 

54 return ( 

55 jsonify( 

56 { 

57 "error": "Scheduler API control is disabled. " 

58 "Contact your administrator to enable it." 

59 } 

60 ), 

61 403, 

62 ) 

63 return f(*args, **kwargs) 

64 

65 return wrapper 

66 

67 

68def safe_error_message(e: Exception, context: str = "") -> str: 

69 """ 

70 Return a safe error message that doesn't expose internal details. 

71 

72 Args: 

73 e: The exception 

74 context: Optional context about what was being attempted 

75 

76 Returns: 

77 A generic error message safe for external users 

78 """ 

79 # Log the actual error for debugging 

80 logger.exception(f"Error in {context}") 

81 

82 # Return generic messages based on exception type 

83 if isinstance(e, ValueError): 

84 return "Invalid input provided" 

85 if isinstance(e, KeyError): 

86 return "Required data missing" 

87 if isinstance(e, TypeError): 

88 return "Invalid data format" 

89 # Generic message for production 

90 return f"An error occurred{f' while {context}' if context else ''}" 

91 

92 

93def _is_job_owned_by_user(job, username, scheduler): 

94 """Check if an APScheduler job belongs to a specific user.""" 

95 # Primary: all news scheduler jobs pass username as first arg 

96 if hasattr(job, "args") and job.args and job.args[0] == username: 

97 return True 

98 # Fallback: check the tracked scheduled_jobs set 

99 if hasattr(scheduler, "user_sessions"): 

100 session_info = scheduler.user_sessions.get(username, {}) 

101 if job.id in session_info.get("scheduled_jobs", set()): 

102 return True 

103 return False 

104 

105 

106def _call_start_research_internal(request_data: dict) -> dict: 

107 """Start a research run by invoking the research route handler in-process. 

108 

109 Both manual subscription runs (``run_subscription_now``) and the overdue 

110 sweep (``check_overdue_subscriptions``) previously issued a loopback HTTP 

111 POST to ``/research/api/start`` via ``safe_post``. That endpoint lives on a 

112 CSRF-protected blueprint (only ``api_v1`` is exempt — see 

113 ``web/app_factory.py``), and a server-to-server request cannot carry a CSRF 

114 token, so every loopback failed with HTTP 400 ("The CSRF token is 

115 missing"). Forwarding the user's session cookie did not help — CSRF is 

116 checked before authentication. The result: both the "run now" button and 

117 the overdue endpoint were broken (only the scheduler path, which calls the 

118 programmatic API directly, still worked). 

119 

120 Calling the view function directly skips the HTTP layer (and therefore the 

121 CSRF ``before_request`` hook) entirely, and removes any need to relay the 

122 session cookie. Must be called from within the authenticated request 

123 context of the caller; that context's session and DB session are 

124 propagated into the nested request context so ``start_research`` can resolve 

125 the user's DB password (it reads ``session["session_id"]`` -> password 

126 store, or ``g.user_password``) and reuse the open connection. 

127 

128 Returns the route's JSON body as a dict (with at least a ``status`` key). 

129 """ 

130 from flask import current_app, g, session 

131 from ..web.routes.research_routes import start_research 

132 from ..database.session_context import get_g_db_session 

133 

134 host_url = request.host_url.rstrip("/") 

135 # Snapshot the caller's auth context. Copying only session["username"] 

136 # would make start_research() -> resolve_user_password() fail with 

137 # "session expired" on encrypted databases, because the DB password is 

138 # keyed by session["session_id"] in the password store. 

139 outer_session = dict(session) 

140 outer_user_password = getattr(g, "user_password", None) 

141 db_session = get_g_db_session() 

142 

143 # Pushing a request context for the same app reuses the existing app 

144 # context, so ``g`` is shared with the caller and ``session`` is fresh. 

145 # session.update() is therefore required; the g.* assignments below are 

146 # usually no-ops but kept so this stays correct if a fresh app context is 

147 # ever pushed (e.g. a different invocation path or a future Flask change). 

148 with current_app.test_request_context( 

149 "/research/api/start", 

150 method="POST", 

151 json=request_data, 

152 base_url=host_url, 

153 ): 

154 session.update(outer_session) 

155 if outer_user_password is not None: 

156 g.user_password = outer_user_password # gitleaks:allow 

157 if db_session is not None: 

158 g.db_session = db_session 

159 

160 result = start_research() 

161 

162 # start_research returns either a Response or (Response, status_code). 

163 # The target path is under /api/, so @login_required returns a JSON 401 

164 # (never an HTML redirect) and every other return is jsonify(...), so 

165 # get_json() always yields a dict — keep this route under /api/. 

166 resp_obj = result[0] if isinstance(result, tuple) else result 

167 data: dict = resp_obj.get_json() 

168 return data 

169 

170 

171# Create Blueprint - no url_prefix here since parent blueprint already has /news 

172news_api_bp = Blueprint("news_api", __name__, url_prefix="/api") 

173# NOTE: Routes use session["username"] (not .get()) intentionally. 

174# @login_required guarantees the key exists; direct access fails fast 

175# if the decorator is ever removed. 

176 

177# Components are initialized in api.py 

178 

179 

180def get_user_id(): 

181 """Get current user ID from session""" 

182 from ..web.auth.decorators import current_user 

183 

184 username = current_user() 

185 

186 if not username: 

187 # For news, we need authenticated users 

188 return None 

189 

190 return username 

191 

192 

193@news_api_bp.route("/feed", methods=["GET"]) 

194@login_required 

195def get_news_feed() -> Any: 

196 """ 

197 Get personalized news feed for user. 

198 

199 Query params: 

200 user_id: User identifier (default: anonymous) 

201 limit: Maximum number of cards to return (default: 20) 

202 use_cache: Whether to use cached news (default: true) 

203 strategy: Override default recommendation strategy 

204 focus: Optional focus area for news 

205 """ 

206 try: 

207 # Get current user (login_required ensures we have one) 

208 user_id = get_user_id() 

209 logger.info(f"News feed requested by user: {user_id}") 

210 

211 # Get query parameters 

212 settings_manager = get_settings_manager() 

213 default_limit = settings_manager.get_setting("news.feed.default_limit") 

214 limit = int(request.args.get("limit", default_limit)) 

215 limit = max(1, min(limit, NEWS_FEED_MAX_LIMIT)) 

216 use_cache = request.args.get("use_cache", "true").lower() == "true" 

217 strategy = request.args.get("strategy") 

218 focus = request.args.get("focus") 

219 subscription_id = request.args.get("subscription_id") 

220 

221 logger.info( 

222 f"News feed params: limit={limit}, subscription_id={subscription_id}, focus={focus}" 

223 ) 

224 

225 # Call the direct API function (now synchronous) 

226 result = api.get_news_feed( 

227 user_id=user_id, 

228 limit=limit, 

229 use_cache=use_cache, 

230 focus=focus, 

231 search_strategy=strategy, 

232 subscription_id=subscription_id, 

233 ) 

234 

235 # Check for errors in result 

236 if "error" in result and result.get("news_items") == []: 

237 # Sanitize error message before returning to client 

238 safe_msg = safe_error_message( 

239 Exception(result["error"]), context="get_news_feed" 

240 ) 

241 return jsonify( 

242 {"error": safe_msg, "news_items": []} 

243 ), 400 if "must be between" in result["error"] else 500 

244 

245 # Debug: Log the result before returning 

246 logger.info( 

247 f"API returning {len(result.get('news_items', []))} news items" 

248 ) 

249 if result.get("news_items"): 

250 logger.info( 

251 f"First item ID: {result['news_items'][0].get('id', 'NO_ID')}" 

252 ) 

253 

254 return jsonify(result) 

255 

256 except Exception as e: 

257 return jsonify( 

258 { 

259 "error": safe_error_message(e, "getting news feed"), 

260 "news_items": [], 

261 } 

262 ), 500 

263 

264 

265@news_api_bp.route("/subscribe", methods=["POST"]) 

266@login_required 

267@require_json_body(error_message="No JSON data provided") 

268def create_subscription() -> Any: 

269 """ 

270 Create a new subscription for user. 

271 

272 JSON body: 

273 query: Search query or topic 

274 subscription_type: "search" or "topic" (default: "search") 

275 refresh_minutes: Refresh interval in minutes (default: from settings) 

276 """ 

277 try: 

278 data = request.get_json(force=True) 

279 except Exception: 

280 # Handle invalid JSON 

281 return jsonify({"error": "Invalid JSON data"}), 400 

282 

283 try: 

284 # Get current user 

285 user_id = get_user_id() 

286 

287 # Extract parameters 

288 query = data.get("query") 

289 subscription_type = data.get("subscription_type", "search") 

290 refresh_minutes = data.get( 

291 "refresh_minutes" 

292 ) # Will use default from api.py 

293 

294 # Extract model configuration (optional) 

295 model_provider = normalize_provider(data.get("model_provider")) 

296 model = data.get("model") 

297 search_strategy = data.get("search_strategy", "news_aggregation") 

298 custom_endpoint = data.get("custom_endpoint") 

299 

300 # Extract additional fields 

301 name = data.get("name") 

302 folder_id = data.get("folder_id") 

303 is_active = data.get("is_active", True) 

304 search_engine = data.get("search_engine") 

305 search_iterations = data.get("search_iterations") 

306 questions_per_iteration = data.get("questions_per_iteration") 

307 

308 # Validate required fields 

309 if not query: 

310 return jsonify({"error": "query is required"}), 400 

311 

312 # Call the direct API function 

313 result = api.create_subscription( 

314 user_id=user_id, 

315 query=query, 

316 subscription_type=subscription_type, 

317 refresh_minutes=refresh_minutes, 

318 model_provider=model_provider, 

319 model=model, 

320 search_strategy=search_strategy, 

321 custom_endpoint=custom_endpoint, 

322 name=name, 

323 folder_id=folder_id, 

324 is_active=is_active, 

325 search_engine=search_engine, 

326 search_iterations=search_iterations, 

327 questions_per_iteration=questions_per_iteration, 

328 ) 

329 

330 return jsonify(result) 

331 

332 except ValueError as e: 

333 return jsonify( 

334 {"error": safe_error_message(e, "creating subscription")} 

335 ), 400 

336 except Exception as e: 

337 return jsonify( 

338 {"error": safe_error_message(e, "creating subscription")} 

339 ), 500 

340 

341 

342@news_api_bp.route("/vote", methods=["POST"]) 

343@login_required 

344@require_json_body(error_message="No JSON data provided") 

345def vote_on_news() -> Any: 

346 """ 

347 Submit vote on a news item. 

348 

349 JSON body: 

350 card_id: ID of the news card 

351 vote: "up" or "down" 

352 """ 

353 try: 

354 data = request.get_json() 

355 

356 # Get current user 

357 user_id = get_user_id() 

358 

359 card_id = data.get("card_id") 

360 vote = data.get("vote") 

361 

362 # Validate 

363 if not all([card_id, vote]): 

364 return jsonify({"error": "card_id and vote are required"}), 400 

365 

366 # Call the direct API function 

367 result = api.submit_feedback( 

368 card_id=card_id, user_id=user_id, vote=vote 

369 ) 

370 

371 return jsonify(result) 

372 

373 except ValueError as e: 

374 error_msg = str(e) 

375 if "not found" in error_msg.lower(): 

376 return jsonify({"error": "Resource not found"}), 404 

377 return jsonify({"error": safe_error_message(e, "submitting vote")}), 400 

378 except Exception as e: 

379 return jsonify({"error": safe_error_message(e, "submitting vote")}), 500 

380 

381 

382@news_api_bp.route("/feedback/batch", methods=["POST"]) 

383@login_required 

384@require_json_body(error_message="No JSON data provided") 

385def get_batch_feedback() -> Any: 

386 """ 

387 Get feedback (votes) for multiple news cards. 

388 JSON body: 

389 card_ids: List of card IDs 

390 """ 

391 try: 

392 data = request.get_json() 

393 card_ids = data.get("card_ids", []) 

394 if not card_ids: 

395 return jsonify({"votes": {}}) 

396 

397 # Get current user 

398 user_id = get_user_id() 

399 

400 # Call the direct API function 

401 result = api.get_votes_for_cards(card_ids=card_ids, user_id=user_id) 

402 

403 return jsonify(result) 

404 

405 except ValueError as e: 

406 error_msg = str(e) 

407 if "not found" in error_msg.lower(): 407 ↛ 409line 407 didn't jump to line 409 because the condition on line 407 was always true

408 return jsonify({"error": "Resource not found"}), 404 

409 return jsonify({"error": safe_error_message(e, "getting votes")}), 400 

410 except Exception as e: 

411 logger.exception("Error getting batch feedback") 

412 return jsonify({"error": safe_error_message(e, "getting votes")}), 500 

413 

414 

415@news_api_bp.route("/feedback/<card_id>", methods=["POST"]) 

416@login_required 

417@require_json_body(error_message="No JSON data provided") 

418def submit_feedback(card_id: str) -> Any: 

419 """ 

420 Submit feedback (vote) for a news card. 

421 

422 JSON body: 

423 vote: "up" or "down" 

424 """ 

425 try: 

426 data = request.get_json() 

427 

428 # Get current user 

429 user_id = get_user_id() 

430 vote = data.get("vote") 

431 

432 # Validate 

433 if not vote: 

434 return jsonify({"error": "vote is required"}), 400 

435 

436 # Call the direct API function 

437 result = api.submit_feedback( 

438 card_id=card_id, user_id=user_id, vote=vote 

439 ) 

440 

441 return jsonify(result) 

442 

443 except ValueError as e: 

444 error_msg = str(e) 

445 if "not found" in error_msg.lower(): 

446 return jsonify({"error": "Resource not found"}), 404 

447 if "must be" in error_msg.lower(): 

448 return jsonify({"error": "Invalid input value"}), 400 

449 return jsonify( 

450 {"error": safe_error_message(e, "submitting feedback")} 

451 ), 400 

452 except Exception as e: 

453 return jsonify( 

454 {"error": safe_error_message(e, "submitting feedback")} 

455 ), 500 

456 

457 

458@news_api_bp.route("/research/<card_id>", methods=["POST"]) 

459@login_required 

460def research_news_item(card_id: str) -> Any: 

461 """ 

462 Perform deeper research on a news item. 

463 

464 JSON body: 

465 depth: "quick", "detailed", or "report" (default: "quick") 

466 """ 

467 try: 

468 data = request.get_json() or {} 

469 depth = data.get("depth", "quick") 

470 

471 # Call the API function which handles the research 

472 result = api.research_news_item(card_id, depth) 

473 

474 return jsonify(result) 

475 

476 except Exception as e: 

477 return jsonify( 

478 {"error": safe_error_message(e, "researching news item")} 

479 ), 500 

480 

481 

482@news_api_bp.route("/subscriptions/current", methods=["GET"]) 

483@login_required 

484def get_current_user_subscriptions() -> Any: 

485 """Get all subscriptions for current user.""" 

486 try: 

487 # Get current user 

488 user_id = get_user_id() 

489 

490 # Ensure we have a database session for the user 

491 # This will trigger register_activity 

492 logger.debug(f"Getting news feed for user {user_id}") 

493 

494 # Use the API function 

495 result = api.get_subscriptions(user_id) 

496 if "error" in result: 

497 logger.error( 

498 f"Error getting subscriptions for user {user_id}: {result['error']}" 

499 ) 

500 return jsonify({"error": "Failed to retrieve subscriptions"}), 500 

501 return jsonify(result) 

502 

503 except Exception as e: 

504 return jsonify( 

505 {"error": safe_error_message(e, "getting subscriptions")} 

506 ), 500 

507 

508 

509@news_api_bp.route("/subscriptions/<subscription_id>", methods=["GET"]) 

510@login_required 

511def get_subscription(subscription_id: str) -> Any: 

512 """Get a single subscription by ID.""" 

513 try: 

514 # Handle null or invalid subscription IDs 

515 if ( 

516 subscription_id == "null" 

517 or subscription_id == "undefined" 

518 or not subscription_id 

519 ): 

520 return jsonify({"error": "Invalid subscription ID"}), 400 

521 

522 # Get the subscription 

523 subscription = api.get_subscription(subscription_id) 

524 

525 if not subscription: 

526 return jsonify({"error": "Subscription not found"}), 404 

527 

528 return jsonify(subscription) 

529 

530 except Exception as e: 

531 return jsonify( 

532 {"error": safe_error_message(e, "getting subscription")} 

533 ), 500 

534 

535 

536@news_api_bp.route("/subscriptions/<subscription_id>", methods=["PUT"]) 

537@login_required 

538@require_json_body(error_message="No JSON data provided") 

539def update_subscription(subscription_id: str) -> Any: 

540 """Update a subscription.""" 

541 try: 

542 data = request.get_json(force=True) 

543 except Exception: 

544 return jsonify({"error": "Invalid JSON data"}), 400 

545 

546 try: 

547 # Prepare update data 

548 update_data = {} 

549 

550 # Map fields from request to storage format 

551 field_mapping = { 

552 "query": "query_or_topic", 

553 "name": "name", 

554 "refresh_minutes": "refresh_interval_minutes", 

555 "is_active": "is_active", 

556 "folder_id": "folder_id", 

557 "model_provider": "model_provider", 

558 "model": "model", 

559 "search_strategy": "search_strategy", 

560 "custom_endpoint": "custom_endpoint", 

561 "search_engine": "search_engine", 

562 "search_iterations": "search_iterations", 

563 "questions_per_iteration": "questions_per_iteration", 

564 } 

565 

566 for request_field, storage_field in field_mapping.items(): 

567 if request_field in data: 

568 update_data[storage_field] = data[request_field] 

569 

570 # Update subscription 

571 result = api.update_subscription(subscription_id, update_data) 

572 

573 if "error" in result: 

574 # Sanitize error message before returning to client 

575 original_error = result["error"] 

576 result["error"] = safe_error_message( 

577 Exception(original_error), "updating subscription" 

578 ) 

579 if "not found" in original_error.lower(): 

580 return jsonify(result), 404 

581 return jsonify(result), 400 

582 

583 return jsonify(result) 

584 

585 except Exception as e: 

586 return jsonify( 

587 {"error": safe_error_message(e, "updating subscription")} 

588 ), 500 

589 

590 

591@news_api_bp.route("/subscriptions/<subscription_id>", methods=["DELETE"]) 

592@login_required 

593def delete_subscription(subscription_id: str) -> Any: 

594 """Delete a subscription.""" 

595 try: 

596 # Call the direct API function 

597 success = api.delete_subscription(subscription_id) 

598 

599 if success: 

600 return jsonify( 

601 { 

602 "status": "success", 

603 "message": f"Subscription {subscription_id} deleted", 

604 } 

605 ) 

606 return jsonify({"error": "Subscription not found"}), 404 

607 

608 except Exception as e: 

609 return jsonify( 

610 {"error": safe_error_message(e, "deleting subscription")} 

611 ), 500 

612 

613 

614@news_api_bp.route("/subscriptions/<subscription_id>/run", methods=["POST"]) 

615@login_required 

616def run_subscription_now(subscription_id: str) -> Any: 

617 """Manually trigger a subscription to run now.""" 

618 try: 

619 from flask import session 

620 from .core.utils import get_local_date_string 

621 from .subscription_runner import ( 

622 advance_refresh_schedule, 

623 build_subscription_request_data, 

624 ) 

625 from ..database.session_context import get_user_db_session 

626 from ..database.models.news import NewsSubscription 

627 from ..settings.manager import SettingsManager 

628 from datetime import datetime, timezone 

629 

630 username = session["username"] 

631 

632 # Load the subscription from the user's database. Reading the ORM row 

633 # directly (rather than the trimmed api.get_subscriptions() dict, which 

634 # drops model_provider/model/search_strategy/search_engine) ensures the 

635 # manual run honors the subscription's saved model config, matching the 

636 # overdue sweep and the scheduler. 

637 # Read the subscription + build the payload, then release the read 

638 # transaction before the blocking POST below. The per-user encrypted DB 

639 # uses deferred isolation, so a SELECT holds a SHARED lock on the 

640 # SQLCipher file until the transaction ends — and the session is 

641 # request-cached, so it is NOT closed on `with` exit. We therefore 

642 # rollback() explicitly to drop the read lock before the (up to 30s) 

643 # HTTP call, then reopen a short session to advance. 

644 with get_user_db_session(username) as db: 

645 sub = ( 

646 db.query(NewsSubscription) 

647 .filter(NewsSubscription.id == subscription_id) 

648 .first() 

649 ) 

650 if not sub: 

651 return jsonify({"error": "Subscription not found"}), 404 

652 

653 subscription_pk = sub.id 

654 # Snapshot next_refresh for the post-POST compare-and-set (below): 

655 # a fast-failing run's failure handler may reset next_refresh on the 

656 # worker thread, and the advance must not clobber that. 

657 prev_next_refresh = sub.next_refresh 

658 settings_manager = SettingsManager(db) 

659 current_date = get_local_date_string(settings_manager) 

660 

661 request_data = build_subscription_request_data( 

662 query_template=sub.query_or_topic, 

663 current_date=current_date, 

664 triggered_by="manual", 

665 subscription_id=sub.id, 

666 model_provider=sub.model_provider, 

667 model=sub.model, 

668 search_strategy=sub.search_strategy, 

669 search_engine=sub.search_engine, 

670 custom_endpoint=sub.custom_endpoint, 

671 title=sub.name, 

672 ) 

673 # End the read transaction so the SHARED lock is released before the 

674 # blocking POST (exiting the `with` does not — the session is 

675 # request-cached and not closed on exit). 

676 db.rollback() 

677 

678 # Start the research in-process. A loopback HTTP POST to 

679 # /research/api/start cannot pass CSRF validation (see 

680 # _call_start_research_internal), so call the route handler directly. 

681 # The read transaction was already released above so the research 

682 # route can reuse the connection without contending for the lock. 

683 data = _call_start_research_internal(request_data) 

684 

685 if data.get("status") in ("success", "queued"): 

686 # Advance the schedule so a subscription that was also overdue 

687 # is not immediately re-run by the scheduler while this run is 

688 # in flight. If the run later fails, the research failure 

689 # handler resets next_refresh so the scheduler retries it. 

690 # Compare-and-set: only advance if next_refresh is unchanged 

691 # since we read it pre-spawn. A fast-failing run can reset it 

692 # (worker thread) before we get here; in that case leave the 

693 # reset in place rather than clobbering it and re-hiding the 

694 # failed subscription for a full interval. 

695 with get_user_db_session(username) as db: 

696 sub = ( 

697 db.query(NewsSubscription) 

698 .filter(NewsSubscription.id == subscription_pk) 

699 .first() 

700 ) 

701 if sub and sub.next_refresh == prev_next_refresh: 

702 advance_refresh_schedule(sub, datetime.now(timezone.utc)) 

703 db.commit() 

704 return jsonify( 

705 { 

706 "status": "success", 

707 "message": "Research started", 

708 "research_id": data.get("research_id"), 

709 "url": f"/progress/{data.get('research_id')}", 

710 } 

711 ) 

712 return jsonify( 

713 { 

714 "error": data.get( 

715 "message", 

716 data.get("error", "Failed to start research"), 

717 ) 

718 } 

719 ), 500 

720 

721 except Exception as e: 

722 return jsonify( 

723 {"error": safe_error_message(e, "running subscription")} 

724 ), 500 

725 

726 

727@news_api_bp.route("/subscriptions/<subscription_id>/history", methods=["GET"]) 

728@login_required 

729def get_subscription_history(subscription_id: str) -> Any: 

730 """Get research history for a subscription.""" 

731 try: 

732 settings_manager = get_settings_manager() 

733 default_limit = settings_manager.get_setting("news.feed.default_limit") 

734 limit = int(request.args.get("limit", default_limit)) 

735 limit = max(1, min(limit, NEWS_FEED_MAX_LIMIT)) 

736 result = api.get_subscription_history(subscription_id, limit) 

737 if "error" in result: 

738 logger.error( 

739 f"Error getting subscription history: {result['error']}" 

740 ) 

741 return jsonify( 

742 { 

743 "error": "Failed to retrieve subscription history", 

744 "history": [], 

745 } 

746 ), 500 

747 return jsonify(result) 

748 except Exception as e: 

749 return jsonify( 

750 {"error": safe_error_message(e, "getting subscription history")} 

751 ), 500 

752 

753 

754@news_api_bp.route("/preferences", methods=["POST"]) 

755@login_required 

756@require_json_body(error_message="No JSON data provided") 

757def save_preferences() -> Any: 

758 """Save user preferences for news.""" 

759 try: 

760 data = request.get_json() 

761 

762 # Get current user 

763 user_id = get_user_id() 

764 preferences = data.get("preferences", {}) 

765 

766 # Call the direct API function 

767 result = api.save_news_preferences(user_id, preferences) 

768 

769 return jsonify(result) 

770 

771 except Exception as e: 

772 return jsonify( 

773 {"error": safe_error_message(e, "saving preferences")} 

774 ), 500 

775 

776 

777@news_api_bp.route("/categories", methods=["GET"]) 

778@login_required 

779def get_categories() -> Any: 

780 """Get news category distribution.""" 

781 try: 

782 # Call the direct API function 

783 result = api.get_news_categories() 

784 

785 return jsonify(result) 

786 

787 except Exception as e: 

788 return jsonify( 

789 {"error": safe_error_message(e, "getting categories")} 

790 ), 500 

791 

792 

793@news_api_bp.route("/scheduler/status", methods=["GET"]) 

794@login_required 

795def get_scheduler_status() -> Any: 

796 """Get activity-based scheduler status.""" 

797 try: 

798 logger.info("Scheduler status endpoint called") 

799 from flask import session 

800 from ..scheduler.background import get_background_job_scheduler 

801 

802 # Get scheduler instance 

803 scheduler = get_background_job_scheduler() 

804 username = session["username"] 

805 show_all = get_env_setting("news.scheduler.allow_api_control", False) 

806 logger.info( 

807 f"Scheduler instance obtained: is_running={scheduler.is_running}" 

808 ) 

809 

810 # Build status manually to avoid potential deadlock 

811 if show_all: 

812 active_users = ( 

813 len(scheduler.user_sessions) 

814 if hasattr(scheduler, "user_sessions") 

815 else 0 

816 ) 

817 else: 

818 active_users = ( 

819 1 

820 if hasattr(scheduler, "user_sessions") 

821 and username in scheduler.user_sessions 

822 else 0 

823 ) 

824 

825 status = { 

826 "scheduler_available": True, # APScheduler is installed and working 

827 "is_running": scheduler.is_running, 

828 "config": scheduler.config.copy() 

829 if hasattr(scheduler, "config") 

830 else {}, 

831 "active_users": active_users, 

832 "total_scheduled_jobs": 0, 

833 } 

834 

835 # Count scheduled jobs 

836 if hasattr(scheduler, "user_sessions"): 836 ↛ 848line 836 didn't jump to line 848 because the condition on line 836 was always true

837 if show_all: 

838 total_jobs = sum( 

839 len(sess.get("scheduled_jobs", set())) 

840 for sess in scheduler.user_sessions.values() 

841 ) 

842 else: 

843 user_session = scheduler.user_sessions.get(username, {}) 

844 total_jobs = len(user_session.get("scheduled_jobs", set())) 

845 status["total_scheduled_jobs"] = total_jobs 

846 

847 # Also count actual APScheduler jobs 

848 if hasattr(scheduler, "scheduler") and scheduler.scheduler: 

849 try: 

850 apscheduler_jobs = scheduler.scheduler.get_jobs() 

851 if not show_all: 

852 apscheduler_jobs = [ 

853 j 

854 for j in apscheduler_jobs 

855 if _is_job_owned_by_user(j, username, scheduler) 

856 ] 

857 status["apscheduler_job_count"] = len(apscheduler_jobs) 

858 status["apscheduler_jobs"] = [ 

859 { 

860 "id": job.id, 

861 "name": job.name, 

862 "next_run": job.next_run_time.isoformat() 

863 if job.next_run_time 

864 else None, 

865 } 

866 for job in apscheduler_jobs[ 

867 :10 

868 ] # Limit to first 10 for display 

869 ] 

870 except Exception: 

871 logger.exception("Error getting APScheduler jobs") 

872 status["apscheduler_job_count"] = 0 

873 

874 logger.info(f"Status built: {list(status.keys())}") 

875 

876 # Add scheduled_jobs field that JS expects 

877 status["scheduled_jobs"] = status.get("total_scheduled_jobs", 0) 

878 

879 logger.info( 

880 f"Returning status: is_running={status.get('is_running')}, active_users={status.get('active_users')}" 

881 ) 

882 return jsonify(status) 

883 

884 except Exception as e: 

885 return jsonify( 

886 {"error": safe_error_message(e, "getting scheduler status")} 

887 ), 500 

888 

889 

890@news_api_bp.route("/scheduler/start", methods=["POST"]) 

891@login_required 

892@scheduler_control_required 

893def start_scheduler() -> Any: 

894 """Start the subscription scheduler.""" 

895 try: 

896 from flask import current_app 

897 from ..scheduler.background import get_background_job_scheduler 

898 

899 # Get scheduler instance 

900 scheduler = get_background_job_scheduler() 

901 

902 if scheduler.is_running: 

903 return jsonify({"message": "Scheduler is already running"}), 200 

904 

905 # Start the scheduler 

906 scheduler.start() 

907 

908 # Update app reference 

909 current_app.background_job_scheduler = scheduler # type: ignore[attr-defined,unused-ignore] 

910 

911 logger.info("News scheduler started via API") 

912 return jsonify( 

913 { 

914 "status": "success", 

915 "message": "Scheduler started", 

916 "active_users": len(scheduler.user_sessions), 

917 } 

918 ) 

919 

920 except Exception as e: 

921 return jsonify( 

922 {"error": safe_error_message(e, "starting scheduler")} 

923 ), 500 

924 

925 

926@news_api_bp.route("/scheduler/stop", methods=["POST"]) 

927@login_required 

928@scheduler_control_required 

929def stop_scheduler() -> Any: 

930 """Stop the subscription scheduler.""" 

931 try: 

932 from flask import current_app 

933 

934 if ( 

935 hasattr(current_app, "background_job_scheduler") 

936 and current_app.background_job_scheduler 

937 ): 

938 scheduler = current_app.background_job_scheduler 

939 if scheduler.is_running: 

940 scheduler.stop() 

941 logger.info("News scheduler stopped via API") 

942 return jsonify( 

943 {"status": "success", "message": "Scheduler stopped"} 

944 ) 

945 return jsonify({"message": "Scheduler is not running"}), 200 

946 return jsonify({"message": "No scheduler instance found"}), 404 

947 

948 except Exception as e: 

949 return jsonify( 

950 {"error": safe_error_message(e, "stopping scheduler")} 

951 ), 500 

952 

953 

954@news_api_bp.route("/scheduler/check-now", methods=["POST"]) 

955@login_required 

956@scheduler_control_required 

957def check_subscriptions_now() -> Any: 

958 """Manually trigger subscription checking.""" 

959 try: 

960 from flask import current_app 

961 

962 if ( 

963 not hasattr(current_app, "background_job_scheduler") 

964 or not current_app.background_job_scheduler 

965 ): 

966 return jsonify({"error": "Scheduler not initialized"}), 503 

967 

968 scheduler = current_app.background_job_scheduler 

969 if not scheduler.is_running: 969 ↛ 973line 969 didn't jump to line 973 because the condition on line 969 was always true

970 return jsonify({"error": "Scheduler is not running"}), 503 

971 

972 # Run the check subscriptions task immediately 

973 scheduler_instance = current_app.background_job_scheduler 

974 

975 # Get count of due subscriptions 

976 from ..database.models import NewsSubscription as BaseSubscription 

977 from datetime import datetime, timedelta, timezone 

978 

979 with get_user_db_session() as session: 

980 now = datetime.now(timezone.utc) 

981 count = ( 

982 session.query(BaseSubscription) 

983 .filter(BaseSubscription.due_filter(now)) 

984 .count() 

985 ) 

986 

987 # Trigger the check asynchronously via APScheduler with app context 

988 username = get_user_id() 

989 if not username: 

990 return jsonify({"error": "No authenticated user"}), 401 

991 

992 scheduler_instance.scheduler.add_job( 

993 func=scheduler_instance._wrap_job( 

994 scheduler_instance._check_user_overdue_subscriptions 

995 ), 

996 args=[username], 

997 trigger="date", 

998 run_date=datetime.now(timezone.utc) + timedelta(seconds=1), 

999 id=f"manual_check_{username}", 

1000 replace_existing=True, 

1001 ) 

1002 

1003 return jsonify( 

1004 { 

1005 "status": "success", 

1006 "message": f"Checking {count} due subscriptions", 

1007 "count": count, 

1008 } 

1009 ) 

1010 

1011 except Exception as e: 

1012 return jsonify( 

1013 {"error": safe_error_message(e, "checking subscriptions")} 

1014 ), 500 

1015 

1016 

1017@news_api_bp.route("/scheduler/cleanup-now", methods=["POST"]) 

1018@login_required 

1019@scheduler_control_required 

1020def trigger_cleanup() -> Any: 

1021 """Manually trigger cleanup job.""" 

1022 try: 

1023 from ..scheduler.background import get_background_job_scheduler 

1024 from datetime import datetime, UTC, timedelta 

1025 

1026 scheduler = get_background_job_scheduler() 

1027 

1028 if not scheduler.is_running: 

1029 return jsonify({"error": "Scheduler is not running"}), 400 

1030 

1031 # Schedule cleanup to run in 1 second 

1032 scheduler.scheduler.add_job( 

1033 scheduler._wrap_job(scheduler._run_cleanup_with_tracking), 

1034 "date", 

1035 run_date=datetime.now(UTC) + timedelta(seconds=1), 

1036 id="manual_cleanup_trigger", 

1037 replace_existing=True, 

1038 ) 

1039 

1040 return jsonify( 

1041 { 

1042 "status": "triggered", 

1043 "message": "Cleanup job will run within seconds", 

1044 } 

1045 ) 

1046 

1047 except Exception as e: 

1048 return jsonify( 

1049 {"error": safe_error_message(e, "triggering cleanup")} 

1050 ), 500 

1051 

1052 

1053@news_api_bp.route("/scheduler/users", methods=["GET"]) 

1054@login_required 

1055def get_active_users() -> Any: 

1056 """Get summary of active user sessions.""" 

1057 try: 

1058 from flask import session 

1059 from ..scheduler.background import get_background_job_scheduler 

1060 

1061 scheduler = get_background_job_scheduler() 

1062 username = session["username"] 

1063 users_summary = scheduler.get_user_sessions_summary() 

1064 

1065 show_all = get_env_setting("news.scheduler.allow_api_control", False) 

1066 if not show_all: 

1067 users_summary = [ 

1068 u for u in users_summary if u.get("user_id") == username 

1069 ] 

1070 

1071 return jsonify( 

1072 {"active_users": len(users_summary), "users": users_summary} 

1073 ) 

1074 

1075 except Exception as e: 

1076 return jsonify( 

1077 {"error": safe_error_message(e, "getting active users")} 

1078 ), 500 

1079 

1080 

1081@news_api_bp.route("/scheduler/stats", methods=["GET"]) 

1082@login_required 

1083def scheduler_stats() -> Any: 

1084 """Get scheduler statistics and state.""" 

1085 try: 

1086 from ..scheduler.background import get_background_job_scheduler 

1087 from flask import session 

1088 

1089 scheduler = get_background_job_scheduler() 

1090 username = session["username"] 

1091 

1092 # Debug info 

1093 debug_info = { 

1094 "current_user": username, 

1095 "scheduler_running": scheduler.is_running, 

1096 "user_sessions": {}, 

1097 "apscheduler_jobs": [], 

1098 } 

1099 

1100 show_all = get_env_setting("news.scheduler.allow_api_control", False) 

1101 

1102 # Get user session info 

1103 if hasattr(scheduler, "user_sessions"): 1103 ↛ 1122line 1103 didn't jump to line 1122 because the condition on line 1103 was always true

1104 for user, session_info in scheduler.user_sessions.items(): 

1105 if not show_all and user != username: 

1106 continue 

1107 debug_info["user_sessions"][user] = { 

1108 "has_password": bool( 

1109 scheduler._credential_store.retrieve(user) 

1110 ), 

1111 "last_activity": session_info.get( 

1112 "last_activity" 

1113 ).isoformat() 

1114 if session_info.get("last_activity") 

1115 else None, 

1116 "scheduled_jobs_count": len( 

1117 session_info.get("scheduled_jobs", set()) 

1118 ), 

1119 } 

1120 

1121 # Get APScheduler jobs 

1122 if hasattr(scheduler, "scheduler") and scheduler.scheduler: 

1123 jobs = scheduler.scheduler.get_jobs() 

1124 if not show_all: 

1125 jobs = [ 

1126 j 

1127 for j in jobs 

1128 if _is_job_owned_by_user(j, username, scheduler) 

1129 ] 

1130 debug_info["apscheduler_jobs"] = [ 

1131 { 

1132 "id": job.id, 

1133 "name": job.name, 

1134 "next_run": job.next_run_time.isoformat() 

1135 if job.next_run_time 

1136 else None, 

1137 "trigger": str(job.trigger), 

1138 } 

1139 for job in jobs 

1140 ] 

1141 

1142 return jsonify(debug_info) 

1143 

1144 except Exception as e: 

1145 return jsonify( 

1146 {"error": safe_error_message(e, "getting scheduler stats")} 

1147 ), 500 

1148 

1149 

1150@news_api_bp.route("/check-overdue", methods=["POST"]) 

1151@login_required 

1152def check_overdue_subscriptions(): 

1153 """Check and run all overdue subscriptions for the current user.""" 

1154 try: 

1155 from flask import session 

1156 from .subscription_runner import ( 

1157 advance_refresh_schedule, 

1158 build_subscription_request_data, 

1159 ) 

1160 from ..database.session_context import get_user_db_session 

1161 from ..database.models.news import NewsSubscription 

1162 from datetime import datetime, UTC 

1163 

1164 username = session["username"] 

1165 

1166 # Get overdue subscriptions 

1167 overdue_count = 0 

1168 results = [] 

1169 with get_user_db_session(username) as db: 

1170 now = datetime.now(UTC) 

1171 overdue_subs = ( 

1172 db.query(NewsSubscription) 

1173 .filter(NewsSubscription.due_filter(now)) 

1174 .all() 

1175 ) 

1176 

1177 logger.info( 

1178 f"Found {len(overdue_subs)} overdue subscriptions for {username}" 

1179 ) 

1180 

1181 # Get timezone-aware current date using settings 

1182 from .core.utils import get_local_date_string 

1183 from ..settings.manager import SettingsManager 

1184 

1185 settings_manager = SettingsManager(db) 

1186 current_date = get_local_date_string(settings_manager) 

1187 

1188 for sub in overdue_subs: 

1189 # Capture identity up front as plain strings. This loop shares 

1190 # one DB session across every start_research call, and 

1191 # start_research's error path does not roll back — so a failed 

1192 # run can leave the session in a PendingRollbackError state. 

1193 # Reading sub.* again in an error branch would then raise (the 

1194 # row was expired by an earlier commit), collapsing the whole 

1195 # sweep. Snapshotting here keeps the error branches session-free. 

1196 sub_id = str(sub.id) 

1197 sub_label = sub.name or sub.query_or_topic[:50] 

1198 try: 

1199 # Run the subscription using the same pattern as run_subscription_now 

1200 logger.info( 

1201 f"Running overdue subscription: {sub_label[:30]}" 

1202 ) 

1203 

1204 # Snapshot for the post-run compare-and-set (see below). 

1205 prev_next_refresh = sub.next_refresh 

1206 

1207 request_data = build_subscription_request_data( 

1208 query_template=sub.query_or_topic, 

1209 current_date=current_date, 

1210 triggered_by="overdue_check", 

1211 subscription_id=sub.id, 

1212 model_provider=sub.model_provider, 

1213 model=sub.model, 

1214 search_strategy=sub.search_strategy, 

1215 search_engine=sub.search_engine, 

1216 custom_endpoint=sub.custom_endpoint, 

1217 title=sub.name, 

1218 ) 

1219 

1220 # Start the research in-process. A loopback HTTP POST to 

1221 # /research/api/start cannot pass CSRF validation (see 

1222 # _call_start_research_internal), so call the route handler 

1223 # directly — this also removes the session-cookie relay 

1224 # that the loopback used to attempt authentication. 

1225 result = _call_start_research_internal(request_data) 

1226 

1227 if result.get("status") in ("success", "queued"): 

1228 overdue_count += 1 

1229 

1230 # Update subscription's last/next refresh times. 

1231 # Compare-and-set: re-read and skip the advance if a 

1232 # fast-failing run already reset next_refresh (worker 

1233 # thread), so we don't clobber the reset and re-hide it. 

1234 db.refresh(sub) 

1235 if sub.next_refresh == prev_next_refresh: 

1236 advance_refresh_schedule(sub, datetime.now(UTC)) 

1237 db.commit() 

1238 

1239 results.append( 

1240 { 

1241 "id": sub_id, 

1242 "name": sub_label, 

1243 "research_id": result.get("research_id"), 

1244 } 

1245 ) 

1246 else: 

1247 # start_research failed and may have left the shared 

1248 # session dirty; reset it so the next subscription runs. 

1249 db.rollback() 

1250 results.append( 

1251 { 

1252 "id": sub_id, 

1253 "name": sub_label, 

1254 # start_research reports failures under 

1255 # "message"; keep "error" as a fallback for 

1256 # any other shape. 

1257 "error": result.get( 

1258 "message", 

1259 result.get( 

1260 "error", "Failed to start research" 

1261 ), 

1262 ), 

1263 } 

1264 ) 

1265 except Exception as e: 

1266 # Recover the shared session (a failed start_research commit 

1267 # can leave it in a PendingRollbackError state) so the 

1268 # remaining overdue subscriptions in the sweep still run. 

1269 db.rollback() 

1270 logger.exception(f"Error running subscription {sub_id}") 

1271 results.append( 

1272 { 

1273 "id": sub_id, 

1274 "name": sub_label, 

1275 "error": safe_error_message( 

1276 e, "running subscription" 

1277 ), 

1278 } 

1279 ) 

1280 

1281 return jsonify( 

1282 { 

1283 "status": "success", 

1284 "overdue_found": len(overdue_subs), 

1285 "started": overdue_count, 

1286 "results": results, 

1287 } 

1288 ) 

1289 

1290 except Exception as e: 

1291 return jsonify( 

1292 {"error": safe_error_message(e, "checking overdue subscriptions")} 

1293 ), 500 

1294 

1295 

1296# Folder and subscription management routes 

1297@news_api_bp.route("/subscription/folders", methods=["GET"]) 

1298@login_required 

1299def get_folders(): 

1300 """Get all folders for the current user""" 

1301 try: 

1302 user_id = get_user_id() 

1303 

1304 with get_user_db_session() as session: 

1305 manager = FolderManager(session) 

1306 folders = manager.get_user_folders(user_id) 

1307 

1308 return jsonify([folder.to_dict() for folder in folders]) 

1309 

1310 except Exception as e: 

1311 return jsonify({"error": safe_error_message(e, "getting folders")}), 500 

1312 

1313 

1314@news_api_bp.route("/subscription/folders", methods=["POST"]) 

1315@login_required 

1316@require_json_body() 

1317def create_folder(): 

1318 """Create a new folder""" 

1319 try: 

1320 data = request.json 

1321 if not data.get("name"): 

1322 return jsonify({"error": "Folder name is required"}), 400 

1323 

1324 with get_user_db_session() as session: 

1325 manager = FolderManager(session) 

1326 

1327 # Check if folder already exists 

1328 existing = ( 

1329 session.query(SubscriptionFolder) 

1330 .filter_by(name=data["name"]) 

1331 .first() 

1332 ) 

1333 if existing: 

1334 return jsonify({"error": "Folder already exists"}), 409 

1335 

1336 folder = manager.create_folder( 

1337 name=data["name"], 

1338 description=data.get("description"), 

1339 ) 

1340 

1341 return jsonify(folder.to_dict()), 201 

1342 

1343 except Exception as e: 

1344 return jsonify({"error": safe_error_message(e, "creating folder")}), 500 

1345 

1346 

1347@news_api_bp.route("/subscription/folders/<folder_id>", methods=["PUT"]) 

1348@login_required 

1349@require_json_body() 

1350def update_folder(folder_id): 

1351 """Update a folder""" 

1352 try: 

1353 data = request.json 

1354 with get_user_db_session() as session: 

1355 manager = FolderManager(session) 

1356 folder = manager.update_folder(folder_id, **data) 

1357 

1358 if not folder: 

1359 return jsonify({"error": "Folder not found"}), 404 

1360 

1361 return jsonify(folder.to_dict()) 

1362 

1363 except Exception as e: 

1364 return jsonify({"error": safe_error_message(e, "updating folder")}), 500 

1365 

1366 

1367@news_api_bp.route("/subscription/folders/<folder_id>", methods=["DELETE"]) 

1368@login_required 

1369def delete_folder(folder_id): 

1370 """Delete a folder""" 

1371 try: 

1372 move_to = request.args.get("move_to") 

1373 

1374 with get_user_db_session() as session: 

1375 manager = FolderManager(session) 

1376 success = manager.delete_folder(folder_id, move_to) 

1377 

1378 if not success: 

1379 return jsonify({"error": "Folder not found"}), 404 

1380 

1381 return jsonify({"status": "deleted"}), 200 

1382 

1383 except Exception as e: 

1384 return jsonify({"error": safe_error_message(e, "deleting folder")}), 500 

1385 

1386 

1387@news_api_bp.route("/subscription/subscriptions/organized", methods=["GET"]) 

1388@login_required 

1389def get_subscriptions_organized(): 

1390 """Get subscriptions organized by folder""" 

1391 try: 

1392 user_id = get_user_id() 

1393 

1394 with get_user_db_session() as session: 

1395 manager = FolderManager(session) 

1396 organized = manager.get_subscriptions_by_folder(user_id) 

1397 

1398 # get_subscriptions_by_folder already returns JSON-friendly dicts 

1399 # ({"folders": [{"folder": {...}, "subscriptions": [...]}, ...], 

1400 # "uncategorized": [...]}). The previous code called .to_dict() on 

1401 # those plain dicts (AttributeError -> HTTP 500). Flatten into the 

1402 # {folder_name: [subscription, ...]} map the subscriptions UI 

1403 # consumes (Object.values(...) for the "all" view, keyed lookup per 

1404 # folder), with ungrouped subscriptions under "uncategorized". 

1405 result = {} 

1406 for entry in organized.get("folders", []): 

1407 folder_name = entry["folder"].get("name") or entry[ 

1408 "folder" 

1409 ].get("id") 

1410 result[folder_name] = entry["subscriptions"] 

1411 # Merge (don't overwrite) so a user folder literally named 

1412 # "uncategorized" doesn't have its subscriptions dropped by the 

1413 # ungrouped bucket. In the normal case this just sets the key. 

1414 result.setdefault("uncategorized", []).extend( 

1415 organized.get("uncategorized", []) 

1416 ) 

1417 

1418 return jsonify(result) 

1419 

1420 except Exception as e: 

1421 return jsonify( 

1422 {"error": safe_error_message(e, "getting organized subscriptions")} 

1423 ), 500 

1424 

1425 

1426@news_api_bp.route( 

1427 "/subscription/subscriptions/<subscription_id>", methods=["PUT"] 

1428) 

1429@login_required 

1430@require_json_body() 

1431def update_subscription_folder(subscription_id): 

1432 """Update a subscription (mainly for folder assignment)""" 

1433 try: 

1434 data = request.json 

1435 logger.info( 

1436 f"Updating subscription {subscription_id} with data: {data}" 

1437 ) 

1438 

1439 with get_user_db_session() as session: 

1440 # Manually handle the update to ensure next_refresh is recalculated 

1441 from ..database.models.news import ( 

1442 NewsSubscription as BaseSubscription, 

1443 ) 

1444 from datetime import datetime, timedelta, timezone 

1445 

1446 sub = ( 

1447 session.query(BaseSubscription) 

1448 .filter_by(id=subscription_id) 

1449 .first() 

1450 ) 

1451 if not sub: 1451 ↛ 1452line 1451 didn't jump to line 1452 because the condition on line 1451 was never true

1452 return jsonify({"error": "Subscription not found"}), 404 

1453 

1454 # `status` is the single source of truth for active/paused (the 

1455 # scheduler keys on it, not the legacy is_active column). Translate 

1456 # an is_active toggle into status and keep both out of the blind 

1457 # setattr loop below -- otherwise a body like {"is_active": false} 

1458 # would flip only the legacy column while status stayed "active" 

1459 # and the scheduler would keep running the subscription. Mirrors 

1460 # api.update_subscription's translation. 

1461 if "is_active" in data: 

1462 sub.status = "active" if data["is_active"] else "paused" 

1463 if "status" in data: 1463 ↛ 1464line 1463 didn't jump to line 1464 because the condition on line 1463 was never true

1464 sub.status = data["status"] 

1465 

1466 # Update remaining fields 

1467 for key, value in data.items(): 

1468 if hasattr(sub, key) and key not in [ 

1469 "id", 

1470 "user_id", 

1471 "created_at", 

1472 "is_active", 

1473 "status", 

1474 ]: 

1475 setattr(sub, key, value) 

1476 

1477 # Recalculate next_refresh if refresh_interval_minutes changed 

1478 if "refresh_interval_minutes" in data: 1478 ↛ 1479line 1478 didn't jump to line 1479 because the condition on line 1478 was never true

1479 new_minutes = data["refresh_interval_minutes"] 

1480 if sub.last_refresh: 

1481 sub.next_refresh = sub.last_refresh + timedelta( 

1482 minutes=new_minutes 

1483 ) 

1484 else: 

1485 sub.next_refresh = datetime.now(timezone.utc) + timedelta( 

1486 minutes=new_minutes 

1487 ) 

1488 logger.info(f"Recalculated next_refresh: {sub.next_refresh}") 

1489 

1490 sub.updated_at = datetime.now(timezone.utc) 

1491 session.commit() 

1492 

1493 # NewsSubscription has no to_dict(); serialize the fields the UI 

1494 # needs explicitly. is_active is derived from status (the source of 

1495 # truth) so the response stays consistent with the toggle above. 

1496 result = { 

1497 "id": sub.id, 

1498 "name": sub.name, 

1499 "status": sub.status, 

1500 "is_active": sub.status == "active", 

1501 "folder_id": sub.folder_id, 

1502 "refresh_interval_minutes": sub.refresh_interval_minutes, 

1503 "next_refresh": sub.next_refresh.isoformat() 

1504 if sub.next_refresh 

1505 else None, 

1506 "last_refresh": sub.last_refresh.isoformat() 

1507 if sub.last_refresh 

1508 else None, 

1509 } 

1510 logger.info( 

1511 f"Updated subscription result: refresh_interval_minutes={result.get('refresh_interval_minutes')}, next_refresh={result.get('next_refresh')}" 

1512 ) 

1513 return jsonify(result) 

1514 

1515 except Exception as e: 

1516 return jsonify( 

1517 {"error": safe_error_message(e, "updating subscription")} 

1518 ), 500 

1519 

1520 

1521@news_api_bp.route("/subscription/stats", methods=["GET"]) 

1522@login_required 

1523def get_subscription_stats(): 

1524 """Get subscription statistics""" 

1525 try: 

1526 user_id = get_user_id() 

1527 

1528 with get_user_db_session() as session: 

1529 manager = FolderManager(session) 

1530 stats = manager.get_subscription_stats(user_id) 

1531 

1532 return jsonify(stats) 

1533 

1534 except Exception as e: 

1535 return jsonify({"error": safe_error_message(e, "getting stats")}), 500 

1536 

1537 

1538# Error handlers 

1539@news_api_bp.errorhandler(400) 

1540def bad_request(e): 

1541 return jsonify({"error": "Bad request"}), 400 

1542 

1543 

1544@news_api_bp.errorhandler(404) 

1545def not_found(e): 

1546 return jsonify({"error": "Resource not found"}), 404 

1547 

1548 

1549@news_api_bp.errorhandler(500) 

1550def internal_error(e): 

1551 return jsonify({"error": "Internal server error"}), 500 

1552 

1553 

1554@news_api_bp.route("/search-history", methods=["GET"]) 

1555@login_required 

1556def get_search_history(): 

1557 """Get search history for current user.""" 

1558 try: 

1559 # Get username from session 

1560 from ..web.auth.decorators import current_user 

1561 

1562 username = current_user() 

1563 if not username: 

1564 # Not authenticated, return empty history 

1565 return jsonify({"search_history": []}) 

1566 

1567 # Get search history from user's encrypted database 

1568 from ..database.session_context import get_user_db_session 

1569 from ..database.models import UserNewsSearchHistory 

1570 

1571 # Get password from Flask g object (set by middleware) 

1572 from flask import g 

1573 

1574 password = getattr(g, "user_password", None) 

1575 

1576 with get_user_db_session(username, password) as db_session: 

1577 history = ( 

1578 db_session.query(UserNewsSearchHistory) 

1579 .order_by(UserNewsSearchHistory.created_at.desc()) 

1580 .limit(20) 

1581 .all() 

1582 ) 

1583 

1584 return jsonify( 

1585 {"search_history": [item.to_dict() for item in history]} 

1586 ) 

1587 

1588 except Exception as e: 

1589 return jsonify( 

1590 {"error": safe_error_message(e, "getting search history")} 

1591 ), 500 

1592 

1593 

1594@news_api_bp.route("/search-history", methods=["POST"]) 

1595@login_required 

1596def add_search_history(): 

1597 """Add a search to the history.""" 

1598 try: 

1599 # Get username from session 

1600 from ..web.auth.decorators import current_user 

1601 

1602 username = current_user() 

1603 if not username: 

1604 # Not authenticated 

1605 return jsonify({"error": "Authentication required"}), 401 

1606 

1607 data = request.get_json() 

1608 logger.info( 

1609 f"add_search_history received data keys: {list(data.keys()) if data else 'None'}" 

1610 ) 

1611 if not data or not data.get("query"): 

1612 logger.warning("Invalid search history data: missing query") 

1613 return jsonify({"error": "query is required"}), 400 

1614 

1615 # Add to user's encrypted database 

1616 from ..database.session_context import get_user_db_session 

1617 from ..database.models import UserNewsSearchHistory 

1618 

1619 # Get password from Flask g object (set by middleware) 

1620 from flask import g 

1621 

1622 password = getattr(g, "user_password", None) 

1623 

1624 with get_user_db_session(username, password) as db_session: 

1625 search_history = UserNewsSearchHistory( 

1626 query=data["query"], 

1627 search_type=data.get("type", "filter"), 

1628 result_count=data.get("resultCount", 0), 

1629 ) 

1630 db_session.add(search_history) 

1631 db_session.commit() 

1632 

1633 return jsonify({"status": "success", "id": search_history.id}) 

1634 

1635 except Exception as e: 

1636 logger.exception("Error adding search history") 

1637 return jsonify( 

1638 {"error": safe_error_message(e, "adding search history")} 

1639 ), 500 

1640 

1641 

1642@news_api_bp.route("/search-history", methods=["DELETE"]) 

1643@login_required 

1644def clear_search_history(): 

1645 """Clear all search history for current user.""" 

1646 try: 

1647 # Get username from session 

1648 from ..web.auth.decorators import current_user 

1649 

1650 username = current_user() 

1651 if not username: 

1652 return jsonify({"status": "success"}) 

1653 

1654 # Clear from user's encrypted database 

1655 from ..database.session_context import get_user_db_session 

1656 from ..database.models import UserNewsSearchHistory 

1657 

1658 # Get password from Flask g object (set by middleware) 

1659 from flask import g 

1660 

1661 password = getattr(g, "user_password", None) 

1662 

1663 with get_user_db_session(username, password) as db_session: 

1664 db_session.query(UserNewsSearchHistory).delete() 

1665 db_session.commit() 

1666 

1667 return jsonify({"status": "success"}) 

1668 

1669 except Exception as e: 

1670 return jsonify( 

1671 {"error": safe_error_message(e, "clearing search history")} 

1672 ), 500