88 FAILED_RETRY_INTERVAL_HOURS ,
99 LIST_UPDATE_INTERVAL_HOURS ,
1010 MAX_CONCURRENT_ONBOARDINGS ,
11+ STUCK_ONBOARDING_LIST_TIMEOUT_HOURS ,
12+ STUCK_RECURRENT_LIST_TIMEOUT_HOURS ,
1113)
1214
1315from .connection import get_db_connection
@@ -50,21 +52,33 @@ async def acquire_list(sql_query: str, params: tuple = None) -> MailingList | No
5052
5153
5254async def acquire_onboarding_list () -> MailingList | None :
53- """Acquire a list that has never been processed (onboarding), bounded by MAX_CONCURRENT_ONBOARDINGS."""
55+ """Acquire a list that has never been processed (onboarding), bounded by
56+ MAX_CONCURRENT_ONBOARDINGS. Also reclaims onboarding lists stuck in PROCESSING
57+ past STUCK_ONBOARDING_LIST_TIMEOUT_HOURS (e.g. worker died mid-clone) — such
58+ rows don't count against the concurrency cap either, so a fully-stuck queue
59+ can't deadlock the cap forever.
60+ """
5461 sql_query = f"""
5562 WITH current_onboarding_count AS (
5663 SELECT COUNT(*) as count
5764 FROM mailinglist."listProcessing" lp
5865 WHERE lp.state = $1
5966 AND lp."lastProcessedAt" IS NULL
67+ AND lp."lockedAt" >= NOW() - INTERVAL '1 hour' * $4::numeric
6068 ),
6169 selected_list AS (
6270 SELECT l.id
6371 FROM mailinglist.lists l
6472 JOIN mailinglist."listProcessing" lp ON lp."listId" = l.id
6573 CROSS JOIN current_onboarding_count c
66- WHERE lp.state = $2
67- AND lp."lockedAt" IS NULL
74+ WHERE (
75+ (lp.state = $2 AND lp."lockedAt" IS NULL)
76+ OR (
77+ lp.state = $1
78+ AND lp."lastProcessedAt" IS NULL
79+ AND lp."lockedAt" < NOW() - INTERVAL '1 hour' * $4::numeric
80+ )
81+ )
6882 AND l."deletedAt" IS NULL
6983 AND c.count < $3
7084 ORDER BY lp.priority ASC, lp."createdAt" ASC
@@ -83,23 +97,40 @@ async def acquire_onboarding_list() -> MailingList | None:
8397 """
8498 return await acquire_list (
8599 sql_query ,
86- (ListState .PROCESSING , ListState .PENDING , MAX_CONCURRENT_ONBOARDINGS ),
100+ (
101+ ListState .PROCESSING ,
102+ ListState .PENDING ,
103+ MAX_CONCURRENT_ONBOARDINGS ,
104+ STUCK_ONBOARDING_LIST_TIMEOUT_HOURS ,
105+ ),
87106 )
88107
89108
90109async def acquire_recurrent_list () -> MailingList | None :
91- """Acquire a previously-processed list that is due for reprocessing."""
110+ """Acquire a previously-processed list that is due for reprocessing. Also
111+ reclaims recurrent lists stuck in PROCESSING past
112+ STUCK_RECURRENT_LIST_TIMEOUT_HOURS (e.g. worker died mid-fetch) regardless of
113+ their normal reprocessing schedule, since a stuck row is by definition overdue.
114+ """
92115 sql_query = f"""
93116 WITH selected_list AS (
94117 SELECT l.id
95118 FROM mailinglist.lists l
96119 JOIN mailinglist."listProcessing" lp ON lp."listId" = l.id
97- WHERE NOT (lp.state = ANY($2))
98- AND lp."lockedAt" IS NULL
99- AND l."deletedAt" IS NULL
100- AND lp."lastProcessedAt" < NOW() - INTERVAL '1 hour' * (
101- CASE WHEN lp.state = $4 THEN $5::numeric ELSE $3::numeric END
120+ WHERE (
121+ (
122+ NOT (lp.state = ANY($2))
123+ AND lp."lastProcessedAt" < NOW() - INTERVAL '1 hour' * (
124+ CASE WHEN lp.state = $4 THEN $5::numeric ELSE $3::numeric END
125+ )
102126 )
127+ OR (
128+ lp.state = $1
129+ AND lp."lastProcessedAt" IS NOT NULL
130+ AND lp."lockedAt" < NOW() - INTERVAL '1 hour' * $6::numeric
131+ )
132+ )
133+ AND l."deletedAt" IS NULL
103134 ORDER BY lp.priority ASC, lp."lastProcessedAt" ASC
104135 LIMIT 1
105136 FOR UPDATE OF lp SKIP LOCKED
@@ -123,6 +154,7 @@ async def acquire_recurrent_list() -> MailingList | None:
123154 LIST_UPDATE_INTERVAL_HOURS ,
124155 ListState .FAILED ,
125156 FAILED_RETRY_INTERVAL_HOURS ,
157+ STUCK_RECURRENT_LIST_TIMEOUT_HOURS ,
126158 ),
127159 )
128160
0 commit comments