Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
24 changes: 17 additions & 7 deletions dtable_events/automations/automations_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -277,17 +277,21 @@ def scan_rules(self):
'per_day_check_time': per_day_check_time,
'per_week_check_time': per_week_check_time,
'per_month_check_time': per_month_check_time
})
}).fetchall()
except Exception as e:
auto_rule_logger.exception('Failed to query scheduled automation rules: %s', e)
db_session.close()
return
finally:
db_session.close()

cached_exceed_keys_set = set()
gen_exceed_key = lambda owner, org_id: org_id if org_id != -1 else owner

try:
for rule in rules:
db_session = None
for rule in rules:

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[Warning] 缺少故障隔离回归测试

Why this matters:
本次修改修复的是 Redis/数据库读取异常导致后续定时规则整体不执行的生产故障,但 PR 未新增覆盖该异常路径的测试。后续调整 session 生命周期或异常处理时,循环隔离与失败后重建 session 的关键行为可能再次退化而不被 CI 发现。

Suggested fix: 为 scan_rules 添加测试,模拟某条规则的 is_exceed 抛出连接异常,断言后续规则仍入队;同时验证下一次查询会使用新建的 session。

try:
if db_session is None:
db_session = self._db_session_class()
automation_task = AutomationTask(
rule_id=rule.id,
run_condition=rule.run_condition,
Expand All @@ -313,9 +317,15 @@ def scan_rules(self):
continue
self.put_task(automation_task)
self.scheduled_trigger_count += 1
except Exception as e:
auto_rule_logger.exception(e)
finally:
except Exception:
auto_rule_logger.exception('run scheduled rule %s error', rule.id)
if db_session is not None:
try:
db_session.close()
except Exception:
pass
db_session = None
if db_session is not None:
db_session.close()

def scheduled_scan(self):
Expand Down
Loading