Hacktoberfest 2026: the issues maintainers tagged for October, open and beginner-friendly. Browse Hacktoberfest issues

fix(worker): 队列消息 reclaim 不查活跃租约——排队中的长任务被「死信→重入队」循环反复冲洗(DLQ 增长 + churn)

Open
#47 0 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

Assessment

Difficulty
4/5
Estimated time
3-5 days
Newbie friendliness
70/100
Issue type
Bug
Clarity
Clearly specified
Activity status
Active
Tech stack
python, redis

Research direction

Start at plaita/server/task_queue.py _reclaim_one() (~line 481) and trace how XCLAIM fires without consulting the lease key {ns}:execution:lease:{execution_id}, then read plaita/server/flow_worker.py _dead_letter_guard() (~line 1914) to see why 'queued, no lease' is treated as 'holder dead'. First run: reproduce queue saturation so a message waits > max_deliveries × claim period and watch times_delivered on plaita:flow:queue:v2 via XPENDING. Done means the four acceptance conditions in the issue hold: no delivery growth/DLQ additions while queued, kill -9 orphan recovery still re-enqueues, short-task throughput unchanged, and XPENDING entries stable in the saturation window.

Written by the indexing model from the issue text.

Description

优先级 P2 · 依赖:无

基线 HEAD = 44febf5

Summary

现象(实测):plaita:flow:queue:v2:dlq 3 条,全部 reason=max_deliveries=5(21:21:45 / 21:28:00 / 21:34:50 生成),payload 全为 #87 的同一条 start 任务(execution_id=99c142cce2b14539aca7a95b0e966391,tenant=default,任务入队时间 19:19:07);三条的 source_id 互为链条(1791371947969 → 1791379305016 → 1791379680771)=「死信→重入队新副本→再烧 5 次投递→再死信」的循环,循环周期 ≈6.4min。同窗口 XPENDING 全部压在 worker-mac-1(主控实测 15 条);21:40 主机 worker 重启后 claim 存量 PEL 又触发一波「交接重入队」,XLEN 15→25。

代码链路(亲读,origin/main):

  • reclaim 只按 idle 判空闲,不查执行租约:plaita/server/task_queue.py:481 _reclaim_one() —— xpending_range 扫描(窗口 8)→ idle < claim_min_idle_ms 跳过 → 否则 XCLAIM 抢走并直接返回该任务重投(return task,deliveries < max 分支);租约/执行状态在本路径零查询。默认 DEFAULT_CLAIM_MIN_IDLE_MS = 60_000(task_queue.py:18)。
  • 死信守卫只覆盖「死信决策」,且把「非终态+无租约」当「持有者已死」:plaita/server/flow_worker.py:1902-1911 接线 dead_letter_guard→_dead_letter_guard()(:1914):非终态 + 租约空 → 重新入队一份同体消息(delivery 归 1)后放行原消息死信——该分支在「派发≠执行、排队等待中(exec 已建但未开跑、无租约)」的常态下会被反复触发(当前队列饱和常态:keeper 超闸 + 双 worker 满载,排队数小时并不罕见)。
  • 即:「排队中(未开跑)」与「持有者崩溃(孤儿)」在守卫眼里不可区分——前者本应静静等槽,实际被判定「死了」→ 死信+重投循环;后者才需要恢复路径。
  • 补充观察:当 exec 真在跑(租约活)时,守卫会跳过死信(行为正确),但其队列消息仍在每个 claim 周期被 XCLAIM 一次 → delivery_count 持续虚增;触顶后每轮「跳过」也是 churn 噪声(消息留 pending 至终态)。

对照租约证据:同一 exec 当前 plaita:execution:lease:99c142cc… 存在且 TTL=86s(活)、status=running(start_time=21:50:04)——即队列里那条被反复死信重投的 start 任务,对应的执行此刻正在正常跑。

影响

  • DLQ 常态增长 + 队列 churn:只要队列饱和(排队时长 > max_deliveries × claim 周期 ≈ 30-35min),排队消息就会被「死信→重入队」循环反复冲洗——DLQ 被无意义条目污染(真死信判别被淹没)、XLEN 虚增(15→25)、worker 重启窗口还会叠加「交接重入队」风暴。
  • 恢复路径可信度受损:真孤儿(持锁 worker 崩溃)与「只是在排队」走同一分支,运维无法从 DLQ 分辨「需要处理的死信」与「排队过久的正常单」。
  • 无双重执行(租约+ack 兜底有效)——本条只针对 churn 与 DLQ 噪声,不是数据正确性问题。

建议

  1. reclaim 前查活跃租约:_reclaim_one 在 XCLAIM 前查 {ns}:execution:lease:{execution_id}(或先 XINFO/状态键):租约活跃 → 跳过该条(不 XCLAIM、不计数),等其终态;仅对「无租约 && 非终态 && 超过更大宽限(如按 start_time 判真正孤儿)」者进入恢复路径。
  2. 死信判据与「排队中」解耦:给「非终态+无租约」分支加排队宽限(例如距 start_time/入队时间 < N 小时一律跳过,只留 pending),或在 exec 状态里区分 queued/started 两态供守卫判据使用(避免用租约缺失代理「持有者已死」)。
  3. 验收条件(可执行):
    • 制造队列饱和(>30min 排队):该消息 times_delivered 不再增长、DLQ 零新增;
    • 真孤儿(可 kill -9 持有 worker):仍能按恢复路径重入队并被接管(不要因修复 #1/#2 而切断恢复);
    • 短任务(<60s 处理)吞吐不受影响;
    • XPENDING 在饱和窗内条目数稳定,不出现「同一条目反复刷新 idle」。

边界

  • 不改 max_deliveries 契约与 DLQ 本身;修复对象是 claim/死信判据。
  • 与 reaper(60m idle + 租约门)分工:reaper 管「执行态终态化」,本条管「队列消息不因排队被误判死信」——两者判据应对齐(都以租约+状态为准)。
Dominant language
Python
Stars
0
Forks
1
Avg merge
1h 44m
Merged PRs (30d)
7

Getting set up

This project ships no dev container, Dockerfile or contributing guide, so setting up is up to you: start from its README, and see our first-contribution guide for the general steps.

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

More from jeffkit/plaita

All issues in jeffkit/plaita

Similar issues

More Python issues

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.