File size: 5,110 Bytes
a37e6db
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
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())