Skip to content

Commit 00c79d6

Browse files
committed
fix: avoid duplicate dispatch within a single rss batch
1 parent 67ef402 commit 00c79d6

4 files changed

Lines changed: 88 additions & 5 deletions

File tree

main.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,7 @@
1515
"astrbot_rss",
1616
"AstrBot RSS Forwarder",
1717
"面向 AstrBot 的 RSS/RSSHub 推送编排插件",
18-
"0.3.2",
18+
"0.3.3",
1919
)
2020
class RSSPlugin(Star, RSSCommands):
2121
def __init__(self, context: Context, config=None):

metadata.yaml

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,6 @@
11
name: astrbot_plugin_rss_forwarder
22
display_name: AstrBot RSS Forwarder
33
desc: 面向 AstrBot 的 RSS/RSSHub 推送编排插件,支持去重、路由与后续智能增强。
4-
version: v0.3.2
4+
version: v0.3.3
55
author: RhoninSeiei
66
repo: https://github.com/RhoninSeiei/astrbot_plugin_rss_forwarder

scheduler.py

Lines changed: 18 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -238,10 +238,12 @@ async def _run_job_once_guarded(self, job: JobConfig) -> None:
238238
pushed_count = 0
239239
parsed_count = 0
240240
skipped_seen_count = 0
241+
skipped_batch_duplicate_count = 0
241242
skipped_history_count = 0
242243
skipped_invalid_target_count = 0
243244
dispatch_fail_count = 0
244245
error_summary = ""
246+
seen_in_run: set[str] = set()
245247

246248
try:
247249
raw_items = await self._call_fetch(job)
@@ -255,9 +257,22 @@ async def _run_job_once_guarded(self, job: JobConfig) -> None:
255257
parsed_count = len(items)
256258
for item in items:
257259
item_id = self._storage.build_dedup_key(item)
258-
if not item_id or await self._storage.has_seen(item_id):
260+
if not item_id:
259261
skipped_seen_count += 1
260262
continue
263+
if item_id in seen_in_run:
264+
skipped_batch_duplicate_count += 1
265+
logger.warning(
266+
"skip job=%s duplicate item in current batch: id=%s title=%s",
267+
job.id,
268+
item_id,
269+
str(item.get("title", "")).strip(),
270+
)
271+
continue
272+
if await self._storage.has_seen(item_id):
273+
skipped_seen_count += 1
274+
continue
275+
seen_in_run.add(item_id)
261276
if self._should_mark_history_only(item, feed_state_map, bootstrap_only=True):
262277
await self._storage.mark_seen(
263278
item_id,
@@ -315,12 +330,13 @@ async def _run_job_once_guarded(self, job: JobConfig) -> None:
315330
error_summary=error_summary,
316331
)
317332
logger.info(
318-
"job=%s finished: fetched=%s parsed=%s pushed=%s skipped_seen=%s skipped_history=%s skipped_invalid_target=%s dispatch_fail=%s duration_ms=%s error=%s",
333+
"job=%s finished: fetched=%s parsed=%s pushed=%s skipped_seen=%s skipped_batch_duplicate=%s skipped_history=%s skipped_invalid_target=%s dispatch_fail=%s duration_ms=%s error=%s",
319334
job.id,
320335
fetched_count,
321336
parsed_count,
322337
pushed_count,
323338
skipped_seen_count,
339+
skipped_batch_duplicate_count,
324340
skipped_history_count,
325341
skipped_invalid_target_count,
326342
dispatch_fail_count,

tests/test_scheduler.py

Lines changed: 68 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,11 @@
77

88
astrbot_module = types.ModuleType("astrbot")
99
astrbot_api_module = types.ModuleType("astrbot.api")
10-
astrbot_api_module.logger = types.SimpleNamespace(info=lambda *a, **k: None)
10+
astrbot_api_module.logger = types.SimpleNamespace(
11+
info=lambda *a, **k: None,
12+
warning=lambda *a, **k: None,
13+
error=lambda *a, **k: None,
14+
)
1115
sys.modules.setdefault("astrbot", astrbot_module)
1216
sys.modules["astrbot.api"] = astrbot_api_module
1317

@@ -146,6 +150,69 @@ async def dispatch(self, item):
146150
self.assertEqual(storage.marked, [("item-1", 123)])
147151

148152

153+
class SchedulerBatchDedupTests(unittest.IsolatedAsyncioTestCase):
154+
async def test_duplicate_items_in_same_batch_are_dispatched_once(self):
155+
class FakeStorage:
156+
def __init__(self):
157+
self.marked = []
158+
159+
def build_dedup_key(self, item):
160+
return item["guid"]
161+
162+
async def has_seen(self, item_id):
163+
return False
164+
165+
async def mark_seen(self, item_id, ttl_seconds=0):
166+
self.marked.append((item_id, ttl_seconds))
167+
168+
async def get_feed_state(self, feed_id):
169+
return {"last_success_time": 0}
170+
171+
async def update_feed_state(self, *args, **kwargs):
172+
return {}
173+
174+
class FakeFetcher:
175+
async def fetch(self, job):
176+
return [{"feed_id": "feed-1"}]
177+
178+
class FakeParser:
179+
def parse(self, raw_items, job):
180+
return [
181+
{"feed_id": "feed-1", "guid": "dup-1", "title": "Same Item", "published_at": ""},
182+
{"feed_id": "feed-1", "guid": "dup-1", "title": "Same Item", "published_at": ""},
183+
]
184+
185+
class FakeDispatcher:
186+
def __init__(self):
187+
self.calls = 0
188+
189+
async def dispatch(self, item):
190+
self.calls += 1
191+
return DispatchResult(success_count=1)
192+
193+
config = types.SimpleNamespace(
194+
jobs=[],
195+
dedup_ttl_seconds=123,
196+
poll_interval_seconds=300,
197+
)
198+
job = types.SimpleNamespace(id="job-1", feed_ids=["feed-1"], enabled=True, interval_seconds=300)
199+
storage = FakeStorage()
200+
dispatcher = FakeDispatcher()
201+
scheduler = RSSScheduler(
202+
config=config,
203+
fetcher=FakeFetcher(),
204+
parser=FakeParser(),
205+
dispatcher=dispatcher,
206+
storage=storage,
207+
pipeline=None,
208+
)
209+
210+
await scheduler._run_job_once_guarded(job)
211+
212+
self.assertEqual(dispatcher.calls, 1)
213+
self.assertEqual(storage.marked, [("dup-1", 123)])
214+
215+
149216
class SchedulerTranslationTest(unittest.IsolatedAsyncioTestCase):
150217
async def test_test_translation_returns_pipeline_error_when_missing(self):
151218
config = types.SimpleNamespace(jobs=[])

0 commit comments

Comments
 (0)