Files
wehub-resource-sync 7a0da7932b
OSV-Scanner (Scheduled) / scan-scheduled (push) Failing after 0s
Create Release / test-gate (push) Has been cancelled
Create Release / release-gate (push) Has been cancelled
Create Release / ci-gate (push) Has been cancelled
Create Release / version-check (push) Has been cancelled
Create Release / e2e-test-gate (push) Has been cancelled
Create Release / responsive-test-gate (push) Has been cancelled
Create Release / compat-test-gate (push) Has been cancelled
Create Release / compose-integration-gate (push) Has been cancelled
Create Release / vulture-gate (push) Has been cancelled
Create Release / build (push) Has been cancelled
Create Release / provenance (push) Has been cancelled
Create Release / prerelease-docker (push) Has been cancelled
Create Release / publish-docker (push) Has been cancelled
Create Release / create-release (push) Has been cancelled
Create Release / cleanup-changelog (push) Has been cancelled
Create Release / trigger-pypi (push) Has been cancelled
Create Release / monitor-pypi (push) Has been cancelled
Create Release / Clean up orphan prerelease tags and signatures (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [research-form] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [research-metrics] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [research-workflow] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [settings-core] (push) Has been cancelled
CodeQL Advanced / Analyze (javascript-typescript) (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [history-news] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [library] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [link-analytics] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [chat-core] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [chat-lifecycle] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [error-benchmark] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [settings-pages] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) (push) Has been cancelled
Docker Tests (Consolidated) / Accessibility Tests (push) Has been cancelled
Docker Tests (Consolidated) / LLM Unit Tests (push) Has been cancelled
Docker Tests (Consolidated) / LLM Example Tests (push) Has been cancelled
Docker Tests (Consolidated) / Production Image Smoke Test (push) Has been cancelled
Docker Tests (Consolidated) / Infrastructure Tests (push) Has been cancelled
OSSF Scorecard / OSSF Security Scorecard Analysis (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [mobile] (push) Has been cancelled
Backwards Compatibility / Verify Encryption Constants (push) Has been cancelled
Backwards Compatibility / PyPI Version Compatibility (push) Has been cancelled
Backwards Compatibility / Database Migration Tests (push) Has been cancelled
CodeQL Advanced / Analyze (python) (push) Has been cancelled
Docker Tests (Consolidated) / detect-changes (push) Has been cancelled
Docker Tests (Consolidated) / Build Test Image (push) Has been cancelled
Docker Tests (Consolidated) / All Pytest Tests + Coverage (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [accessibility] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [api-crud] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [auth-login] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [auth-pages] (push) Has been cancelled
Docker Tests (Consolidated) / UI Tests (Puppeteer) [auth-register] (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:08:55 +08:00

2458 lines
85 KiB
Python

"""
Tests for news/subscription_manager/scheduler.py
Tests cover:
- BackgroundJobScheduler singleton pattern
- Configuration loading
- User session management
- Scheduler lifecycle
- update_user_info method
- unregister_user method
- _schedule_user_subscriptions method
- _schedule_document_processing method
- _get_document_scheduler_settings method
- invalidate_user_settings_cache method
- invalidate_all_settings_cache method
- _check_subscription method
- trigger_document_processing method
- get_document_scheduler_status method
"""
import pytest
from unittest.mock import Mock, patch, MagicMock
import threading
from datetime import datetime, timedelta, UTC
from apscheduler.jobstores.base import JobLookupError
class TestNewsSchedulerSingleton:
"""Tests for BackgroundJobScheduler singleton pattern."""
def test_news_scheduler_is_singleton(self):
"""BackgroundJobScheduler follows singleton pattern."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
# Reset singleton for test
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
scheduler1 = BackgroundJobScheduler()
scheduler2 = BackgroundJobScheduler()
assert scheduler1 is scheduler2
def test_scheduler_has_required_attributes(self):
"""BackgroundJobScheduler has required attributes after init."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
# Reset singleton for test
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
scheduler = BackgroundJobScheduler()
assert hasattr(scheduler, "user_sessions")
assert hasattr(scheduler, "lock")
assert hasattr(scheduler, "scheduler")
assert hasattr(scheduler, "config")
assert hasattr(scheduler, "is_running")
class TestSchedulerConfiguration:
"""Tests for scheduler configuration."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_default_config_values(self, scheduler):
"""Default configuration has expected values."""
config = scheduler.config
assert config["enabled"] is True
assert config["retention_hours"] == 48
assert config["cleanup_interval_hours"] == 1
assert config["max_jitter_seconds"] == 300
assert config["max_concurrent_jobs"] == 10
assert config["subscription_batch_size"] == 5
assert config["activity_check_interval_minutes"] == 5
def test_initialize_with_settings(self, scheduler):
"""Scheduler can be initialized with settings manager."""
mock_settings = Mock()
mock_settings.get.return_value = None
# Should not raise
scheduler.initialize_with_settings(mock_settings)
assert scheduler.settings_manager is mock_settings
class TestSchedulerLifecycle:
"""Tests for scheduler start/stop lifecycle."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler_instance = MagicMock()
mock_scheduler.return_value = mock_scheduler_instance
instance = BackgroundJobScheduler()
yield instance
def test_scheduler_initial_state_not_running(self, scheduler):
"""Scheduler is not running initially."""
assert scheduler.is_running is False
def test_user_sessions_initially_empty(self, scheduler):
"""User sessions dict is initially empty."""
assert scheduler.user_sessions == {}
class TestUserSessionManagement:
"""Tests for user session tracking."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_lock_is_thread_lock(self, scheduler):
"""Scheduler has threading lock for thread safety."""
assert isinstance(scheduler.lock, type(threading.Lock()))
class TestSchedulerAvailability:
"""Tests for scheduler availability flag."""
def test_scheduler_is_available(self):
"""Scheduler availability flag is True."""
from local_deep_research.scheduler.background import (
SCHEDULER_AVAILABLE,
)
assert SCHEDULER_AVAILABLE is True
class TestSchedulerStart:
"""Tests for scheduler start method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler_instance = MagicMock()
mock_scheduler.return_value = mock_scheduler_instance
instance = BackgroundJobScheduler()
instance.set_app(MagicMock())
yield instance
def test_start_sets_is_running(self, scheduler):
"""Starting scheduler sets is_running to True."""
scheduler.start()
assert scheduler.is_running is True
scheduler.scheduler.start.assert_called_once()
def test_start_when_disabled(self, scheduler):
"""Scheduler doesn't start when disabled."""
scheduler.config["enabled"] = False
scheduler.start()
assert scheduler.is_running is False
scheduler.scheduler.start.assert_not_called()
def test_start_when_already_running(self, scheduler):
"""Scheduler warns when already running."""
scheduler.is_running = True
scheduler.start()
# Should not call start again
scheduler.scheduler.start.assert_not_called()
def test_start_adds_cleanup_job(self, scheduler):
"""Starting scheduler adds cleanup job."""
scheduler.start()
# Check that add_job was called at least once
assert scheduler.scheduler.add_job.called
class TestSchedulerStop:
"""Tests for scheduler stop method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler_instance = MagicMock()
mock_scheduler.return_value = mock_scheduler_instance
instance = BackgroundJobScheduler()
yield instance
def test_stop_sets_is_running_false(self, scheduler):
"""Stopping scheduler sets is_running to False."""
scheduler.is_running = True
scheduler.stop()
assert scheduler.is_running is False
def test_stop_when_not_running(self, scheduler):
"""Stopping scheduler when not running is safe."""
scheduler.is_running = False
# Should not raise
scheduler.stop()
assert scheduler.is_running is False
class TestGetSetting:
"""Tests for _get_setting method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_get_setting_with_settings_manager(self, scheduler):
"""_get_setting uses settings manager when available."""
mock_settings = Mock()
mock_settings.get_setting.return_value = 100
scheduler.settings_manager = mock_settings
result = scheduler._get_setting("some.key", 50)
assert result == 100
mock_settings.get_setting.assert_called_once_with(
"some.key", default=50
)
def test_get_setting_without_settings_manager(self, scheduler):
"""_get_setting returns default without settings manager."""
# No settings manager
result = scheduler._get_setting("some.key", 50)
assert result == 50
class TestSchedulerStatus:
"""Tests for scheduler status methods."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_get_status_when_not_running(self, scheduler):
"""Get status when scheduler is not running."""
scheduler.is_running = False
if hasattr(scheduler, "get_status"):
status = scheduler.get_status()
assert (
status.get("running") is False
or status.get("is_running") is False
)
def test_get_status_when_running(self, scheduler):
"""Get status when scheduler is running."""
scheduler.is_running = True
if hasattr(scheduler, "get_status"):
status = scheduler.get_status()
assert (
status.get("running") is True
or status.get("is_running") is True
)
class TestSchedulerRegisterUser:
"""Tests for user registration methods."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_register_user_adds_to_sessions(self, scheduler):
"""Registering user adds them to sessions dict."""
if hasattr(scheduler, "register_user_activity"):
scheduler.register_user_activity("testuser", "password123")
assert "testuser" in scheduler.user_sessions
def test_register_user_updates_activity(self, scheduler):
"""Registering existing user updates last_activity."""
if hasattr(scheduler, "register_user_activity"):
scheduler.register_user_activity("testuser", "password123")
first_activity = scheduler.user_sessions["testuser"].get(
"last_activity"
)
# Register again
import time
time.sleep(0.1)
scheduler.register_user_activity("testuser", "password123")
second_activity = scheduler.user_sessions["testuser"].get(
"last_activity"
)
# Activity should be updated
if first_activity and second_activity:
assert second_activity >= first_activity
class TestSchedulerUnregisterUser:
"""Tests for user unregistration."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_unregister_removes_user(self, scheduler):
"""Unregistering user removes them from sessions."""
# Set up proper session structure
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": [],
"last_activity": None,
}
scheduler._credential_store.store("testuser", "test")
if hasattr(scheduler, "unregister_user"):
scheduler.unregister_user("testuser")
assert "testuser" not in scheduler.user_sessions
def test_unregister_nonexistent_user(self, scheduler):
"""Unregistering non-existent user is safe."""
if hasattr(scheduler, "unregister_user"):
# Should not raise
scheduler.unregister_user("nonexistent")
class TestScheduleUserSubscriptions:
"""Tests for _schedule_user_subscriptions method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_schedule_user_subscriptions_uses_jitter(self, scheduler):
"""_schedule_user_subscriptions applies random jitter."""
# Verify the scheduler has max_jitter_seconds config
assert "max_jitter_seconds" in scheduler.config
assert scheduler.config["max_jitter_seconds"] == 300
def test_schedule_user_subscriptions_respects_batch_size(self, scheduler):
"""_schedule_user_subscriptions respects subscription_batch_size."""
assert "subscription_batch_size" in scheduler.config
assert scheduler.config["subscription_batch_size"] == 5
def test_schedule_user_subscriptions_jitter_calculation(self, scheduler):
"""Jitter is calculated based on max_jitter_seconds."""
import random
random.seed(42) # Make deterministic for test
max_jitter = scheduler.config["max_jitter_seconds"]
# Generate some jitter values
jitters = [random.randint(0, max_jitter) for _ in range(10)]
# All values should be within range
assert all(0 <= j <= max_jitter for j in jitters)
def test_schedule_user_subscriptions_schedules_jobs(self, scheduler):
"""_schedule_user_subscriptions adds jobs to the scheduler."""
if hasattr(scheduler, "_schedule_user_subscriptions"):
# Method exists
assert callable(scheduler._schedule_user_subscriptions)
class TestProcessUserDocuments:
"""Tests for _process_user_documents method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_process_user_documents_batch_processing(self, scheduler):
"""_process_user_documents processes in batches."""
# Verify batch size config exists
assert "subscription_batch_size" in scheduler.config
def test_process_user_documents_max_concurrent(self, scheduler):
"""_process_user_documents respects max_concurrent_jobs."""
assert "max_concurrent_jobs" in scheduler.config
assert scheduler.config["max_concurrent_jobs"] == 10
class TestStoreResearchResult:
"""Tests for _store_research_result method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_store_research_result_serialization(self, scheduler):
"""Research results are properly serialized."""
# The scheduler should have retention_hours configured
assert "retention_hours" in scheduler.config
assert scheduler.config["retention_hours"] == 48
class TestCleanupOldResults:
"""Tests for cleanup functionality."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_cleanup_interval_configured(self, scheduler):
"""Cleanup interval is properly configured."""
assert "cleanup_interval_hours" in scheduler.config
assert scheduler.config["cleanup_interval_hours"] == 1
def test_retention_hours_configured(self, scheduler):
"""Retention hours is properly configured."""
assert "retention_hours" in scheduler.config
assert scheduler.config["retention_hours"] == 48
class TestActivityTracking:
"""Tests for user activity tracking."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_activity_check_interval_configured(self, scheduler):
"""Activity check interval is properly configured."""
assert "activity_check_interval_minutes" in scheduler.config
assert scheduler.config["activity_check_interval_minutes"] == 5
def test_inactive_user_detection(self, scheduler):
"""Inactive users can be detected."""
from datetime import datetime, timedelta, UTC
if hasattr(scheduler, "user_sessions"):
# Set up a user session with old activity
old_activity = datetime.now(UTC) - timedelta(hours=1)
scheduler.user_sessions["old_user"] = {
"scheduled_jobs": [],
"last_activity": old_activity,
}
scheduler._credential_store.store("old_user", "test")
# The user session should be in the dict
assert "old_user" in scheduler.user_sessions
class TestSchedulerExceptionHandling:
"""Tests for scheduler exception handling."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_scheduler_handles_job_exceptions(self, scheduler):
"""Scheduler handles exceptions in job execution."""
# The scheduler should have proper error handling
assert scheduler.scheduler is not None
def test_scheduler_recovers_from_errors(self, scheduler):
"""Scheduler can recover from errors."""
scheduler.is_running = True
# Stopping should work even after errors
scheduler.stop()
assert scheduler.is_running is False
# =============================================================================
# Phase 2 Tests: Comprehensive testing of critical scheduler methods
# =============================================================================
class TestUpdateUserInfo:
"""Tests for update_user_info method - CRITICAL for user session management."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_creates_new_session_for_new_user(self, scheduler):
"""New user gets session entry created."""
scheduler.is_running = True
with patch.object(
scheduler, "_schedule_user_subscriptions"
) as mock_schedule:
scheduler.update_user_info("newuser", "password123")
assert "newuser" in scheduler.user_sessions
mock_schedule.assert_called_once_with("newuser")
def test_stores_password_in_credential_store(self, scheduler):
"""Password is stored in credential store, not in session dict."""
scheduler.is_running = True
with patch.object(scheduler, "_schedule_user_subscriptions"):
scheduler.update_user_info("testuser", "mySecretPassword")
assert "password" not in scheduler.user_sessions["testuser"]
assert (
scheduler._credential_store.retrieve("testuser")
== "mySecretPassword"
)
def test_sets_last_activity_time(self, scheduler):
"""Activity timestamp is set on new user."""
scheduler.is_running = True
with patch.object(scheduler, "_schedule_user_subscriptions"):
before = datetime.now(UTC)
scheduler.update_user_info("testuser", "password")
after = datetime.now(UTC)
last_activity = scheduler.user_sessions["testuser"]["last_activity"]
assert before <= last_activity <= after
def test_initializes_empty_scheduled_jobs_set(self, scheduler):
"""Jobs set is empty initially for new user."""
scheduler.is_running = True
with patch.object(scheduler, "_schedule_user_subscriptions"):
scheduler.update_user_info("testuser", "password")
assert (
scheduler.user_sessions["testuser"]["scheduled_jobs"] == set()
)
def test_updates_existing_user_password(self, scheduler):
"""Password update works for existing user."""
scheduler.is_running = True
scheduler.user_sessions["existinguser"] = {
"last_activity": datetime.now(UTC) - timedelta(hours=1),
"scheduled_jobs": {"job1"},
}
scheduler._credential_store.store("existinguser", "oldpassword")
with patch.object(scheduler, "_schedule_user_subscriptions"):
scheduler.update_user_info("existinguser", "newpassword")
assert (
scheduler._credential_store.retrieve("existinguser")
== "newpassword"
)
# Jobs should be preserved
assert (
"job1"
in scheduler.user_sessions["existinguser"]["scheduled_jobs"]
)
def test_updates_last_activity_for_existing_user(self, scheduler):
"""Activity timestamp is updated for existing user."""
scheduler.is_running = True
old_time = datetime.now(UTC) - timedelta(hours=1)
scheduler.user_sessions["existinguser"] = {
"last_activity": old_time,
"scheduled_jobs": set(),
}
scheduler._credential_store.store("existinguser", "password")
with patch.object(scheduler, "_schedule_user_subscriptions"):
scheduler.update_user_info("existinguser", "password")
new_activity = scheduler.user_sessions["existinguser"][
"last_activity"
]
assert new_activity > old_time
def test_does_nothing_when_scheduler_not_running(self, scheduler):
"""Graceful return when scheduler is not running."""
scheduler.is_running = False
scheduler.update_user_info("testuser", "password")
assert "testuser" not in scheduler.user_sessions
def test_calls_schedule_user_subscriptions(self, scheduler):
"""Triggers subscription scheduling after user info update."""
scheduler.is_running = True
with patch.object(
scheduler, "_schedule_user_subscriptions"
) as mock_schedule:
scheduler.update_user_info("user1", "pass1")
scheduler.update_user_info("user1", "pass2") # Update existing
assert mock_schedule.call_count == 2
mock_schedule.assert_called_with("user1")
class TestUnregisterUserComprehensive:
"""Comprehensive tests for unregister_user method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_removes_user_from_sessions(self, scheduler):
"""Session is deleted when user unregisters."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
scheduler.unregister_user("testuser")
assert "testuser" not in scheduler.user_sessions
def test_removes_all_scheduled_jobs(self, scheduler):
"""All jobs are cleaned up when user unregisters."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": {"job1", "job2", "job3"},
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
scheduler.unregister_user("testuser")
# Verify remove_job was called for each job
assert scheduler.scheduler.remove_job.call_count == 3
scheduler.scheduler.remove_job.assert_any_call("job1")
scheduler.scheduler.remove_job.assert_any_call("job2")
scheduler.scheduler.remove_job.assert_any_call("job3")
def test_handles_job_lookup_error(self, scheduler):
"""Graceful handling when job not found during removal."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": {"missing_job"},
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
scheduler.scheduler.remove_job.side_effect = JobLookupError(
"missing_job"
)
# Should not raise
scheduler.unregister_user("testuser")
assert "testuser" not in scheduler.user_sessions
def test_invalidates_settings_cache(self, scheduler):
"""Cache is cleared when user unregisters."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
# Add to cache
scheduler._settings_cache["testuser"] = MagicMock()
scheduler.unregister_user("testuser")
assert "testuser" not in scheduler._settings_cache
def test_handles_nonexistent_user(self, scheduler):
"""No error for unknown user."""
# Should not raise
scheduler.unregister_user("nonexistent")
assert scheduler.scheduler.remove_job.call_count == 0
def test_thread_safe_removal(self, scheduler):
"""Lock is used properly during removal."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
# Track if lock was acquired
original_lock = scheduler.lock
lock_acquired = []
class TrackingLock:
def __enter__(self):
lock_acquired.append(True)
return original_lock.__enter__()
def __exit__(self, *args):
return original_lock.__exit__(*args)
scheduler.lock = TrackingLock()
scheduler.unregister_user("testuser")
assert len(lock_acquired) == 1
class TestScheduleUserSubscriptionsComprehensive:
"""Comprehensive tests for _schedule_user_subscriptions - CRITICAL scheduling logic."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_returns_early_if_no_session(self, scheduler):
"""Graceful return without session info."""
# No session set up
with patch.object(
scheduler, "_schedule_document_processing"
) as mock_doc:
scheduler._schedule_user_subscriptions("nonexistent")
# Should not crash, document processing should not be called
mock_doc.assert_not_called()
def test_queries_active_subscriptions(self, scheduler):
"""Database query filters for active subscriptions."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_query = MagicMock()
mock_db.query.return_value = mock_query
mock_query.filter.return_value.all.return_value = []
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
# status is the source of truth for "active" (see
# NewsSubscription.active_filter); the query now uses .filter()
# with that predicate rather than filter_by(is_active=True).
# Compare the actual SQLAlchemy expression so the test pins the
# predicate itself, not merely that .filter() was invoked.
from local_deep_research.database.models.news import (
NewsSubscription,
)
(predicate,) = mock_query.filter.call_args.args
assert predicate.compare(NewsSubscription.active_filter())
def test_clears_old_jobs_before_scheduling(self, scheduler):
"""Old jobs are removed before scheduling new ones."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": {"old_job_1", "old_job_2"},
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_db.query.return_value.filter.return_value.all.return_value = []
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
# Old jobs should have been removed
assert scheduler.scheduler.remove_job.call_count >= 2
def test_calculates_jitter_within_bounds(self, scheduler):
"""Jitter respects max_jitter_seconds config."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
scheduler.config["max_jitter_seconds"] = 100
mock_db = MagicMock()
# Create mock subscription
mock_sub = MagicMock()
mock_sub.id = 1
mock_sub.refresh_interval_minutes = 30
mock_sub.name = "Test Sub"
mock_sub.query_or_topic = "test query"
mock_db.query.return_value.filter.return_value.all.return_value = [
mock_sub
]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch(
"local_deep_research.scheduler.background.random"
) as mock_random:
mock_random.randint.return_value = 50
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
# Verify randint was called with correct bounds
mock_random.randint.assert_called_with(0, 100)
def test_uses_interval_trigger_for_hourly(self, scheduler):
"""<=60 min refresh interval uses interval trigger."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.id = 1
mock_sub.refresh_interval_minutes = 60 # Exactly 60 minutes
mock_sub.name = "Hourly Sub"
mock_sub.query_or_topic = "hourly query"
mock_db.query.return_value.filter.return_value.all.return_value = [
mock_sub
]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
# Should use interval trigger
call_args = scheduler.scheduler.add_job.call_args
assert call_args.kwargs.get("trigger") == "interval"
def test_uses_date_trigger_for_infrequent(self, scheduler):
""">60 min refresh interval uses date trigger."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.id = 1
mock_sub.refresh_interval_minutes = 120 # 2 hours
mock_sub.name = "Infrequent Sub"
mock_sub.query_or_topic = "infrequent query"
mock_sub.next_refresh = None # No previous refresh
mock_db.query.return_value.filter.return_value.all.return_value = [
mock_sub
]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
call_args = scheduler.scheduler.add_job.call_args
assert call_args.kwargs.get("trigger") == "date"
def test_schedules_subscription_with_no_next_refresh(self, scheduler):
"""Subscription with no next_refresh gets scheduled in future."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.id = 1
mock_sub.refresh_interval_minutes = 120 # > 60 min uses date trigger
mock_sub.name = "No Refresh Sub"
mock_sub.query_or_topic = "query"
mock_sub.next_refresh = None # No previous refresh time
mock_db.query.return_value.filter.return_value.all.return_value = [
mock_sub
]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
# Should schedule with date trigger
call_args = scheduler.scheduler.add_job.call_args
assert call_args is not None, "add_job was not called"
assert call_args.kwargs.get("trigger") == "date"
# run_date should be in the future (refresh_minutes + jitter)
run_date = call_args.kwargs.get("run_date")
assert run_date is not None
time_diff = (run_date - datetime.now(UTC)).total_seconds()
# Should be approximately 120 minutes (7200 seconds) + jitter
assert (
7100
<= time_diff
<= 7200 + scheduler.config["max_jitter_seconds"] + 1
)
def test_adds_job_to_scheduler(self, scheduler):
"""Job is registered with scheduler."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.id = 42
mock_sub.refresh_interval_minutes = 30
mock_sub.name = "Test Sub"
mock_sub.query_or_topic = "test query"
mock_db.query.return_value.filter.return_value.all.return_value = [
mock_sub
]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
# Verify add_job was called with correct ID
scheduler.scheduler.add_job.assert_called()
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs["id"] == "testuser_42"
def test_tracks_job_in_session_jobs_set(self, scheduler):
"""Job ID is added to user's scheduled_jobs set."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.id = 99
mock_sub.refresh_interval_minutes = 30
mock_sub.name = "Tracked Sub"
mock_sub.query_or_topic = "tracked query"
mock_db.query.return_value.filter.return_value.all.return_value = [
mock_sub
]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
assert (
"testuser_99"
in scheduler.user_sessions["testuser"]["scheduled_jobs"]
)
def test_handles_database_error(self, scheduler):
"""Exception during database access is handled."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.side_effect = Exception(
"Database connection failed"
)
# Should not raise
with patch.object(scheduler, "_schedule_document_processing"):
scheduler._schedule_user_subscriptions("testuser")
class TestScheduleDocumentProcessing:
"""Tests for _schedule_document_processing method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_returns_early_if_no_session(self, scheduler):
"""Graceful return without session info."""
# No session
scheduler._schedule_document_processing("nonexistent")
# Should not add any jobs
scheduler.scheduler.add_job.assert_not_called()
def test_skips_if_disabled_in_settings(self, scheduler):
"""Respects enabled=False setting."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
disabled_settings = DocumentSchedulerSettings(enabled=False)
with patch.object(
scheduler,
"_get_document_scheduler_settings",
return_value=disabled_settings,
):
scheduler._schedule_document_processing("testuser")
# Should not add any jobs
scheduler.scheduler.add_job.assert_not_called()
def test_removes_existing_document_job(self, scheduler):
"""Old document job is cleaned up."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": {"testuser_document_processing"},
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(
enabled=True, interval_seconds=1800
)
scheduler.scheduler.get_job.return_value = MagicMock()
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
scheduler._schedule_document_processing("testuser")
# The document-processing job is torn down before being re-added.
# (_schedule_document_processing also tears down the opt-in
# library-sweep job, so use assert_any_call rather than asserting the
# last call.)
scheduler.scheduler.remove_job.assert_any_call(
"testuser_document_processing"
)
def test_creates_job_with_correct_interval(self, scheduler):
"""Interval from settings is used."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(
enabled=True, interval_seconds=3600
)
scheduler.scheduler.get_job.return_value = MagicMock()
scheduler.scheduler.remove_job.side_effect = JobLookupError("not found")
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
scheduler._schedule_document_processing("testuser")
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs["seconds"] == 3600
def test_job_has_jitter(self, scheduler):
"""Jitter is applied to prevent simultaneous processing."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(enabled=True)
scheduler.scheduler.get_job.return_value = MagicMock()
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
scheduler._schedule_document_processing("testuser")
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs.get("jitter") == 30
def test_job_has_max_instances_1(self, scheduler):
"""Prevents overlapping document processing."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(enabled=True)
scheduler.scheduler.get_job.return_value = MagicMock()
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
scheduler._schedule_document_processing("testuser")
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs.get("max_instances") == 1
def test_verifies_job_was_added(self, scheduler):
"""Job verification check occurs."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(enabled=True)
mock_job = MagicMock()
mock_job.next_run_time = datetime.now(UTC)
scheduler.scheduler.get_job.return_value = mock_job
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
scheduler._schedule_document_processing("testuser")
# Verify get_job was called to verify the job exists
scheduler.scheduler.get_job.assert_called_with(
"testuser_document_processing"
)
def test_handles_job_lookup_error_on_remove(self, scheduler):
"""Graceful handling when existing job not found."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(enabled=True)
scheduler.scheduler.remove_job.side_effect = JobLookupError("not found")
scheduler.scheduler.get_job.return_value = MagicMock()
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
# Should not raise
scheduler._schedule_document_processing("testuser")
class TestGetDocumentSchedulerSettings:
"""Tests for _get_document_scheduler_settings method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_returns_cached_settings(self, scheduler):
"""Cache hit returns cached value."""
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
cached_settings = DocumentSchedulerSettings(
enabled=True, interval_seconds=999, download_pdfs=True
)
scheduler._settings_cache["cacheduser"] = cached_settings
result = scheduler._get_document_scheduler_settings("cacheduser")
assert result is cached_settings
assert result.interval_seconds == 999
def test_fetches_from_db_on_cache_miss(self, scheduler):
"""DB query on cache miss."""
scheduler.user_sessions["dbuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("dbuser", "password")
mock_db = MagicMock()
mock_sm = MagicMock()
mock_sm.get_setting.side_effect = lambda key, default: {
"document_scheduler.enabled": True,
"document_scheduler.interval_seconds": 2400,
"document_scheduler.download_pdfs": True,
"document_scheduler.extract_text": False,
"document_scheduler.generate_rag": True,
"document_scheduler.last_run": "2024-01-01T00:00:00",
}.get(key, default)
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch(
"local_deep_research.settings.manager.SettingsManager",
return_value=mock_sm,
):
result = scheduler._get_document_scheduler_settings("dbuser")
assert result.enabled is True
assert result.interval_seconds == 2400
assert result.download_pdfs is True
assert result.extract_text is False
assert result.generate_rag is True
def test_caches_fetched_settings(self, scheduler):
"""Result is cached after DB fetch."""
scheduler.user_sessions["cacheuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("cacheuser", "password")
mock_db = MagicMock()
mock_sm = MagicMock()
mock_sm.get_setting.return_value = True
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch(
"local_deep_research.settings.manager.SettingsManager",
return_value=mock_sm,
):
scheduler._get_document_scheduler_settings("cacheuser")
assert "cacheuser" in scheduler._settings_cache
def test_returns_defaults_on_no_session(self, scheduler):
"""Graceful return with defaults without session."""
result = scheduler._get_document_scheduler_settings("nosession")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
defaults = DocumentSchedulerSettings.defaults()
assert result.enabled == defaults.enabled
assert result.interval_seconds == defaults.interval_seconds
def test_returns_defaults_on_db_error(self, scheduler):
"""Error handling returns defaults."""
scheduler.user_sessions["erroruser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("erroruser", "password")
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.side_effect = Exception("DB connection failed")
result = scheduler._get_document_scheduler_settings("erroruser")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
defaults = DocumentSchedulerSettings.defaults()
assert result.enabled == defaults.enabled
def test_force_refresh_bypasses_cache(self, scheduler):
"""Force refresh fetches from DB even with cached value."""
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
cached_settings = DocumentSchedulerSettings(interval_seconds=999)
scheduler._settings_cache["refreshuser"] = cached_settings
scheduler.user_sessions["refreshuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("refreshuser", "password")
mock_db = MagicMock()
mock_sm = MagicMock()
mock_sm.get_setting.side_effect = lambda key, default: {
"document_scheduler.interval_seconds": 5000,
}.get(key, default)
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch(
"local_deep_research.settings.manager.SettingsManager",
return_value=mock_sm,
):
result = scheduler._get_document_scheduler_settings(
"refreshuser", force_refresh=True
)
# Should get fresh value, not cached
assert result.interval_seconds == 5000
def test_settings_are_frozen_dataclass(self, scheduler):
"""Immutability for thread safety verified."""
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(
enabled=True, interval_seconds=1800
)
# Frozen dataclass should raise FrozenInstanceError on modification
with pytest.raises(Exception): # dataclasses.FrozenInstanceError
settings.enabled = False
class TestInvalidateUserSettingsCache:
"""Tests for invalidate_user_settings_cache method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_removes_user_from_cache(self, scheduler):
"""Entry is deleted from cache."""
scheduler._settings_cache["testuser"] = MagicMock()
scheduler.invalidate_user_settings_cache("testuser")
assert "testuser" not in scheduler._settings_cache
def test_returns_true_if_found(self, scheduler):
"""True returned on successful removal."""
scheduler._settings_cache["existinguser"] = MagicMock()
result = scheduler.invalidate_user_settings_cache("existinguser")
assert result is True
def test_returns_false_if_not_found(self, scheduler):
"""False returned when user not in cache."""
result = scheduler.invalidate_user_settings_cache("nonexistent")
assert result is False
def test_thread_safe(self, scheduler):
"""Uses lock for thread safety."""
scheduler._settings_cache["testuser"] = MagicMock()
# Track if lock was acquired
original_lock = scheduler._settings_cache_lock
lock_acquired = []
class TrackingLock:
def __enter__(self):
lock_acquired.append(True)
return original_lock.__enter__()
def __exit__(self, *args):
return original_lock.__exit__(*args)
scheduler._settings_cache_lock = TrackingLock()
scheduler.invalidate_user_settings_cache("testuser")
assert len(lock_acquired) == 1
class TestInvalidateAllSettingsCache:
"""Tests for invalidate_all_settings_cache method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_clears_all_entries(self, scheduler):
"""Cache is emptied."""
scheduler._settings_cache["user1"] = MagicMock()
scheduler._settings_cache["user2"] = MagicMock()
scheduler._settings_cache["user3"] = MagicMock()
scheduler.invalidate_all_settings_cache()
assert len(scheduler._settings_cache) == 0
def test_returns_count_cleared(self, scheduler):
"""Returns correct count of cleared entries."""
scheduler._settings_cache["user1"] = MagicMock()
scheduler._settings_cache["user2"] = MagicMock()
result = scheduler.invalidate_all_settings_cache()
assert result == 2
def test_thread_safe(self, scheduler):
"""Uses lock for thread safety."""
scheduler._settings_cache["testuser"] = MagicMock()
original_lock = scheduler._settings_cache_lock
lock_acquired = []
class TrackingLock:
def __enter__(self):
lock_acquired.append(True)
return original_lock.__enter__()
def __exit__(self, *args):
return original_lock.__exit__(*args)
scheduler._settings_cache_lock = TrackingLock()
scheduler.invalidate_all_settings_cache()
assert len(lock_acquired) == 1
class TestCheckSubscription:
"""Tests for _check_subscription method - CRITICAL subscription refresh logic."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_removes_job_if_no_session(self, scheduler):
"""Cleanup when user gone."""
# No session
scheduler._check_subscription("goneuser", 123)
scheduler.scheduler.remove_job.assert_called_with("goneuser_123")
def test_handles_job_removal_error(self, scheduler):
"""Graceful handling when job removal fails."""
scheduler.scheduler.remove_job.side_effect = JobLookupError("not found")
# Should not raise
scheduler._check_subscription("goneuser", 123)
def test_skips_inactive_subscription(self, scheduler):
"""Respects is_active=False."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.status = "paused"
mock_db.query.return_value.get.return_value = mock_sub
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(
scheduler, "_trigger_subscription_research_sync"
) as mock_trigger:
scheduler._check_subscription("testuser", 1)
mock_trigger.assert_not_called()
def test_replaces_date_placeholder(self, scheduler):
"""YYYY-MM-DD is replaced with actual date."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.status = "active"
mock_sub.query_or_topic = "News from YYYY-MM-DD"
mock_sub.id = 1
mock_sub.name = "Test"
mock_sub.refresh_interval_minutes = 60
mock_sub.model_provider = "openai"
mock_sub.model = "gpt-4"
mock_sub.search_strategy = "news"
mock_sub.search_engine = "google"
mock_db.query.return_value.get.return_value = mock_sub
# Get the job and set up mock
scheduler.scheduler.get_job.return_value = None
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch(
"local_deep_research.news.core.utils.get_local_date_string",
return_value="2024-06-15",
):
with patch.object(
scheduler, "_trigger_subscription_research_sync"
) as mock_trigger:
scheduler._check_subscription("testuser", 1)
call_args = mock_trigger.call_args
subscription_data = call_args[0][1]
assert subscription_data["query"] == "News from 2024-06-15"
assert (
subscription_data["original_query"]
== "News from YYYY-MM-DD"
)
def test_updates_last_refresh_time(self, scheduler):
"""Timestamp is updated after check."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.status = "active"
mock_sub.query_or_topic = "Test query"
mock_sub.id = 1
mock_sub.name = "Test"
mock_sub.refresh_interval_minutes = 60
mock_sub.last_refresh = None
mock_sub.next_refresh = None
mock_db.query.return_value.get.return_value = mock_sub
scheduler.scheduler.get_job.return_value = None
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_trigger_subscription_research_sync"):
scheduler._check_subscription("testuser", 1)
# last_refresh should have been updated
assert mock_sub.last_refresh is not None
mock_db.commit.assert_called()
def test_calculates_next_refresh(self, scheduler):
"""Next run time is calculated."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.status = "active"
mock_sub.query_or_topic = "Test query"
mock_sub.id = 1
mock_sub.name = "Test"
mock_sub.refresh_interval_minutes = 60
mock_db.query.return_value.get.return_value = mock_sub
scheduler.scheduler.get_job.return_value = None
before = datetime.now(UTC)
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_trigger_subscription_research_sync"):
scheduler._check_subscription("testuser", 1)
# next_refresh should be approximately 60 minutes from now
assert mock_sub.next_refresh is not None
time_diff = (mock_sub.next_refresh - before).total_seconds()
# Should be approximately 60 minutes (3600 seconds)
assert 3590 <= time_diff <= 3610
def test_triggers_research(self, scheduler):
"""Research API is called."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.status = "active"
mock_sub.query_or_topic = "Test query"
mock_sub.id = 1
mock_sub.name = "Test Sub"
mock_sub.refresh_interval_minutes = 60
mock_sub.model_provider = "openai"
mock_sub.model = "gpt-4"
mock_sub.search_strategy = "news"
mock_sub.search_engine = "google"
mock_db.query.return_value.get.return_value = mock_sub
scheduler.scheduler.get_job.return_value = None
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(
scheduler, "_trigger_subscription_research_sync"
) as mock_trigger:
scheduler._check_subscription("testuser", 1)
mock_trigger.assert_called_once()
call_args = mock_trigger.call_args[0]
assert call_args[0] == "testuser"
assert call_args[1]["id"] == 1
def test_reschedules_for_date_trigger(self, scheduler):
"""Continues scheduling with date trigger."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.status = "active"
mock_sub.query_or_topic = "Test query"
mock_sub.id = 1
mock_sub.name = "Test"
mock_sub.refresh_interval_minutes = 120 # > 60, so uses date trigger
mock_db.query.return_value.get.return_value = mock_sub
# Mock existing job with DateTrigger
mock_job = MagicMock()
mock_job.trigger.__class__.__name__ = "DateTrigger"
scheduler.scheduler.get_job.return_value = mock_job
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
with patch.object(scheduler, "_trigger_subscription_research_sync"):
scheduler._check_subscription("testuser", 1)
# Should reschedule with add_job
scheduler.scheduler.add_job.assert_called()
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs["trigger"] == "date"
assert call_kwargs["id"] == "testuser_1"
def test_handles_database_error(self, scheduler):
"""Exception is handled gracefully."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.side_effect = Exception("DB error")
# Should not raise
scheduler._check_subscription("testuser", 1)
class TestTriggerDocumentProcessing:
"""Tests for trigger_document_processing method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_returns_false_if_no_session(self, scheduler):
"""Graceful return without session."""
scheduler.is_running = True
result = scheduler.trigger_document_processing("nonexistent")
assert result is False
def test_returns_false_if_not_running(self, scheduler):
"""Graceful return when scheduler stopped."""
scheduler.is_running = False
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
result = scheduler.trigger_document_processing("testuser")
assert result is False
def test_schedules_immediate_job(self, scheduler):
"""Job is scheduled for 1 second from now."""
scheduler.is_running = True
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_job = MagicMock()
mock_job.next_run_time = datetime.now(UTC)
scheduler.scheduler.get_job.return_value = mock_job
before = datetime.now(UTC)
result = scheduler.trigger_document_processing("testuser")
assert result is True
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs["trigger"] == "date"
run_date = call_kwargs["run_date"]
# Should be approximately 1 second from now
time_diff = (run_date - before).total_seconds()
assert 0 <= time_diff <= 2
def test_verifies_job_added(self, scheduler):
"""Verification check occurs."""
scheduler.is_running = True
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_job = MagicMock()
scheduler.scheduler.get_job.return_value = mock_job
scheduler.trigger_document_processing("testuser")
scheduler.scheduler.get_job.assert_called_with(
"testuser_document_processing_manual"
)
def test_returns_false_if_job_not_added(self, scheduler):
"""Returns False when job verification fails."""
scheduler.is_running = True
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
scheduler.scheduler.get_job.return_value = None
result = scheduler.trigger_document_processing("testuser")
assert result is False
def test_handles_exception(self, scheduler):
"""Error is handled gracefully."""
scheduler.is_running = True
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
scheduler.scheduler.add_job.side_effect = Exception("Scheduler error")
result = scheduler.trigger_document_processing("testuser")
assert result is False
class TestGetDocumentSchedulerStatus:
"""Tests for get_document_scheduler_status method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_returns_disabled_for_unknown_user(self, scheduler):
"""Unknown user returns disabled status."""
result = scheduler.get_document_scheduler_status("unknown")
assert result["enabled"] is False
assert "message" in result
def test_includes_all_processing_options(self, scheduler):
"""All processing flags are present."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(
enabled=True,
download_pdfs=True,
extract_text=False,
generate_rag=True,
)
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
result = scheduler.get_document_scheduler_status("testuser")
assert "processing_options" in result
assert result["processing_options"]["download_pdfs"] is True
assert result["processing_options"]["extract_text"] is False
assert result["processing_options"]["generate_rag"] is True
def test_shows_has_scheduled_job(self, scheduler):
"""Job tracking is correct."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": {"testuser_document_processing"},
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(enabled=True)
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
result = scheduler.get_document_scheduler_status("testuser")
assert result["has_scheduled_job"] is True
def test_shows_user_active_status(self, scheduler):
"""Active flag is correct."""
scheduler.user_sessions["activeuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("activeuser", "password")
from local_deep_research.scheduler.background import (
DocumentSchedulerSettings,
)
settings = DocumentSchedulerSettings(enabled=True)
with patch.object(
scheduler, "_get_document_scheduler_settings", return_value=settings
):
result = scheduler.get_document_scheduler_status("activeuser")
assert result["user_active"] is True
def test_handles_exception(self, scheduler):
"""Error returns safe dict."""
scheduler.user_sessions["erroruser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("erroruser", "password")
with patch.object(
scheduler,
"_get_document_scheduler_settings",
side_effect=Exception("Settings error"),
):
result = scheduler.get_document_scheduler_status("erroruser")
assert result["enabled"] is False
assert "message" in result
class TestCheckUserOverdueSubscriptions:
"""Tests for _check_user_overdue_subscriptions method."""
@pytest.fixture
def scheduler(self):
"""Create a fresh scheduler instance."""
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
BackgroundJobScheduler._instance = None
with patch(
"local_deep_research.scheduler.background.BackgroundScheduler"
) as mock_scheduler:
mock_scheduler.return_value = MagicMock()
instance = BackgroundJobScheduler()
yield instance
def test_returns_early_if_no_session(self, scheduler):
"""Graceful return without session."""
# No session
scheduler._check_user_overdue_subscriptions("nonexistent")
# Should not attempt to query database
scheduler.scheduler.add_job.assert_not_called()
def test_finds_overdue_subscriptions(self, scheduler):
"""Queries for overdue subscriptions correctly."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_query = MagicMock()
mock_db.query.return_value = mock_query
mock_query.filter.return_value.all.return_value = []
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
scheduler._check_user_overdue_subscriptions("testuser")
# Verify query was made
mock_db.query.assert_called()
def test_schedules_overdue_with_delay(self, scheduler):
"""Overdue subscriptions are scheduled with random delay."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
mock_db = MagicMock()
mock_sub = MagicMock()
mock_sub.id = 1
mock_sub.name = "Overdue Sub"
mock_sub.query_or_topic = "overdue query"
mock_query = MagicMock()
mock_db.query.return_value = mock_query
mock_query.filter.return_value.all.return_value = [mock_sub]
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.return_value.__enter__ = MagicMock(
return_value=mock_db
)
mock_db_session.return_value.__exit__ = MagicMock(
return_value=False
)
scheduler._check_user_overdue_subscriptions("testuser")
# Should schedule job
scheduler.scheduler.add_job.assert_called()
call_kwargs = scheduler.scheduler.add_job.call_args.kwargs
assert call_kwargs["trigger"] == "date"
def test_handles_database_error(self, scheduler):
"""Exception is handled gracefully."""
scheduler.user_sessions["testuser"] = {
"scheduled_jobs": set(),
"last_activity": datetime.now(UTC),
}
scheduler._credential_store.store("testuser", "password")
with patch(
"local_deep_research.database.session_context.get_user_db_session"
) as mock_db_session:
mock_db_session.side_effect = Exception("DB error")
# Should not raise
scheduler._check_user_overdue_subscriptions("testuser")
class TestSchedulerEgressBackstop:
"""The document scheduler runs on an APScheduler worker thread with no
egress context; _arm_egress_backstop must set one from the user's saved
settings so DownloadService fetches get the audit-hook secondary net
(R2-8). @thread_cleanup clears it on exit."""
def test_arm_egress_backstop_sets_context_from_settings(self):
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
from local_deep_research.security.egress.policy import EgressScope
sched = BackgroundJobScheduler()
fake_sm = MagicMock()
fake_sm.get_settings_snapshot.return_value = {
"policy.egress_scope": "private_only",
"search.tool": "library",
}
fake_sm.get_setting.side_effect = lambda k, d=None: (
"library" if k == "search.tool" else d
)
with patch(
"local_deep_research.security.egress.audit_hook.set_active_context"
) as mock_set:
sched._arm_egress_backstop(fake_sm, "alice")
assert mock_set.call_count == 1
ctx = mock_set.call_args[0][0]
assert ctx.scope == EgressScope.PRIVATE_ONLY
def test_arm_egress_backstop_never_raises_on_bad_settings(self):
from local_deep_research.scheduler.background import (
BackgroundJobScheduler,
)
sched = BackgroundJobScheduler()
fake_sm = MagicMock()
fake_sm.get_settings_snapshot.side_effect = RuntimeError("no db")
# Must swallow the error (best-effort backstop).
sched._arm_egress_backstop(fake_sm, "alice")