File size: 3,183 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
"""
Queue Processing Worker

Background worker that processes supervised execution queue.
Runs every 60 seconds to check for available users and process their queues.
"""

import asyncio
import logging
from datetime import datetime

from core.database import get_db
from core.supervised_queue_service import SupervisedQueueService

logger = logging.getLogger(__name__)


class QueueProcessingWorker:
    """
    Background worker for processing supervised execution queue.

    Runs every 60 seconds to:
    1. Find users who recently became online/away
    2. Fetch their pending queues (batch of 10)
    3. Execute queued entries with supervision
    4. Handle expired queues (>24 hours)

    Performance target: <5s per batch of 10 entries
    """

    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("QueueProcessingWorker started")

        while self.running:
            try:
                start_time = datetime.now()

                # Process pending queues
                await self.process_pending_queues()

                # Mark expired queues
                await self.mark_expired_queues()

                elapsed = (datetime.now() - start_time).total_seconds()
                logger.info(
                    f"QueueProcessingWorker cycle completed in {elapsed:.2f}s"
                )

            except Exception as e:
                logger.error(f"QueueProcessingWorker 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("QueueProcessingWorker stopped")

    async def process_pending_queues(self):
        """Process pending queue entries for available users."""
        db = next(get_db())

        try:
            service = SupervisedQueueService(db)
            processed = await service.process_pending_queues(limit=10)

            if processed:
                logger.info(
                    f"Processed {len(processed)} supervised queue entries"
                )

        finally:
            db.close()

    async def mark_expired_queues(self):
        """Mark expired queue entries as failed."""
        db = next(get_db())

        try:
            service = SupervisedQueueService(db)
            count = await service.mark_expired_queues()

            if count > 0:
                logger.info(f"Marked {count} expired queues as failed")

        finally:
            db.close()


# ============================================================================
# Worker Entry Point
# ============================================================================

async def main():
    """Entry point for running the worker."""
    worker = QueueProcessingWorker(interval_seconds=60)

    try:
        await worker.run()
    except KeyboardInterrupt:
        logger.info("Received keyboard interrupt")
        await worker.stop()


if __name__ == "__main__":
    asyncio.run(main())