Spaces:
Sleeping
Sleeping
File size: 5,110 Bytes
862d728 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 | """
Activity State Worker
Background worker that processes user activity state transitions.
Runs every 60 seconds to check for inactive users and update their states.
State Transitions:
- online → away after 5 minutes of inactivity
- away → offline after 15 minutes of inactivity
- offline → online on activity resumption
"""
import asyncio
import logging
from datetime import datetime
from core.database import get_db
from core.user_activity_service import UserActivityService
logger = logging.getLogger(__name__)
class ActivityStateWorker:
"""
Background worker for processing user activity state transitions.
Runs every 60 seconds to:
1. Check user state transitions based on inactivity
2. Process manual override expiry
3. Clean up stale sessions
4. Emit state change events (optional, for WebSocket notifications)
Performance target: <1s per batch of 100 users
"""
def __init__(self, interval_seconds: int = 60):
self.interval_seconds = interval_seconds
self.running = False
async def run(self):
"""Main worker loop."""
self.running = True
logger.info("ActivityStateWorker started")
while self.running:
try:
start_time = datetime.now()
# Process state transitions
await self.process_state_transitions()
# Process manual override expiry
await self.process_manual_override_expiry()
# Cleanup stale sessions
await self.cleanup_stale_sessions()
elapsed = (datetime.now() - start_time).total_seconds()
logger.info(
f"ActivityStateWorker cycle completed in {elapsed:.2f}s"
)
except Exception as e:
logger.error(f"ActivityStateWorker error: {e}", exc_info=True)
# Wait for next cycle
await asyncio.sleep(self.interval_seconds)
async def stop(self):
"""Stop the worker."""
self.running = False
logger.info("ActivityStateWorker stopped")
async def process_state_transitions(self):
"""Process state transitions for inactive users."""
db = next(get_db())
try:
service = UserActivityService(db)
transitions = await service.transition_state_batch(limit=100)
if transitions["total_processed"] > 0:
transition_summary = ", ".join([
f"{k}: {v}"
for k, v in transitions.items()
if k != "total_processed" and v > 0
])
logger.info(
f"State transitions: {transitions['total_processed']} processed"
+ (f" ({transition_summary})" if transition_summary else "")
)
finally:
db.close()
async def process_manual_override_expiry(self):
"""Process manual overrides that have expired."""
from core.models import UserActivity
from sqlalchemy.orm import Session
db: Session = next(get_db())
try:
# Find expired manual overrides
now = datetime.utcnow()
expired = db.query(UserActivity).filter(
UserActivity.manual_override == True,
UserActivity.manual_override_expires_at.isnot(None),
UserActivity.manual_override_expires_at < now
).all()
count = 0
for activity in expired:
# Clear manual override
activity.manual_override = False
activity.manual_override_expires_at = None
# Recalculate state based on actual activity
service = UserActivityService(db)
await service._recalculate_activity_state(activity)
count += 1
logger.info(
f"Cleared expired manual override for user {activity.user_id}"
)
if count > 0:
db.commit()
logger.info(f"Cleared {count} expired manual overrides")
finally:
db.close()
async def cleanup_stale_sessions(self):
"""Clean up stale sessions (no heartbeat for >1 hour)."""
db = next(get_db())
try:
service = UserActivityService(db)
count = await service.cleanup_stale_sessions(limit=50)
if count > 0:
logger.info(f"Cleaned up {count} stale sessions")
finally:
db.close()
# ============================================================================
# Worker Entry Point
# ============================================================================
async def main():
"""Entry point for running the worker."""
worker = ActivityStateWorker(interval_seconds=60)
try:
await worker.run()
except KeyboardInterrupt:
logger.info("Received keyboard interrupt")
await worker.stop()
if __name__ == "__main__":
asyncio.run(main())
|