fix(blog) : acquitter les GUID après livraison XMPP #39
@@ -66,6 +66,12 @@ processus. Le mode `PRONOTE_AUTH_MODE=qr_token` est incompatible avec cette gara
|
|||||||
le refuse avant toute connexion afin de ne pas désynchroniser le token local du token distant.
|
le refuse avant toute connexion afin de ne pas désynchroniser le token local du token distant.
|
||||||
Le dry-run ne remplace pas une vérification des paramètres réellement chargés.
|
Le dry-run ne remplace pas une vérification des paramètres réellement chargés.
|
||||||
|
|
||||||
|
Si le blog RSS est activé, ses GUID ne sont acquittés qu'après confirmation de
|
||||||
|
l'envoi XMPP. Un refus, une exception, l'absence de canal ou un `--dry-run`
|
||||||
|
laisse donc les articles récupérables à l'exécution suivante ; les en-têtes
|
||||||
|
HTTP associés à ces articles suivent la même règle pour éviter un `304` qui
|
||||||
|
masquerait une livraison non confirmée.
|
||||||
|
|
||||||
En mode `PRONOTE_AUTH_MODE=qr_token`, le fichier
|
En mode `PRONOTE_AUTH_MODE=qr_token`, le fichier
|
||||||
`.pronote_auth_state.json` et son verrou frère sont créés dans le répertoire
|
`.pronote_auth_state.json` et son verrou frère sont créés dans le répertoire
|
||||||
de travail du service (par exemple `/var/lib/pronote-sync`) avec le mode
|
de travail du service (par exemple `/var/lib/pronote-sync`) avec le mode
|
||||||
|
|||||||
@@ -230,11 +230,13 @@ class PipelineRunner:
|
|||||||
data = normalize_step(fetched, generated_at=now)
|
data = normalize_step(fetched, generated_at=now)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
blog_articles = fetch_blog_step(self._blog_client, self._blog_state)
|
blog_result = fetch_blog_step(self._blog_client, self._blog_state)
|
||||||
|
blog_articles = list(blog_result.articles)
|
||||||
except PipelineCriticalError:
|
except PipelineCriticalError:
|
||||||
raise
|
raise
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
self._warn("fetch_blog", self._redact(exc))
|
self._warn("fetch_blog", self._redact(exc))
|
||||||
|
blog_result = None
|
||||||
blog_articles = []
|
blog_articles = []
|
||||||
|
|
||||||
try:
|
try:
|
||||||
@@ -290,8 +292,11 @@ class PipelineRunner:
|
|||||||
)
|
)
|
||||||
if self._channel is not None and not self._dry_run:
|
if self._channel is not None and not self._dry_run:
|
||||||
try:
|
try:
|
||||||
if not send_step(self._channel, message):
|
delivered = send_step(self._channel, message)
|
||||||
|
if not delivered:
|
||||||
self._warn("send", "Le canal XMPP a refusé l'envoi")
|
self._warn("send", "Le canal XMPP a refusé l'envoi")
|
||||||
|
elif blog_result is not None and self._blog_state is not None:
|
||||||
|
self._blog_state.acknowledge(blog_result)
|
||||||
except PipelineCriticalError:
|
except PipelineCriticalError:
|
||||||
raise
|
raise
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
|
|||||||
@@ -2,23 +2,28 @@
|
|||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
from pronote_sync.models.blog import BlogArticle
|
from pronote_sync.sources.blog.result import BlogRSSFetchResult
|
||||||
from pronote_sync.sources.blog.rss import BlogRSSClient
|
from pronote_sync.sources.blog.rss import BlogRSSClient
|
||||||
from pronote_sync.sources.blog.state import BlogRSSState
|
from pronote_sync.sources.blog.state import BlogRSSState
|
||||||
from pronote_sync.utils.redaction import redact_exception
|
from pronote_sync.utils.redaction import redact_exception
|
||||||
|
|
||||||
|
|
||||||
def fetch_blog_step(client: BlogRSSClient | None, state: BlogRSSState | None) -> list[BlogArticle]:
|
def fetch_blog_step(client: BlogRSSClient | None, state: BlogRSSState | None) -> BlogRSSFetchResult:
|
||||||
"""Récupère les articles RSS nouveaux en conservant l'état du client.
|
"""Récupère les articles RSS nouveaux sans les acquitter.
|
||||||
|
|
||||||
|
L'état des GUID est acquitté séparément par le pipeline après confirmation
|
||||||
|
de la livraison XMPP. Les en-têtes de cache d'une réponse sans article
|
||||||
|
peuvent être conservés immédiatement, car aucune livraison n'est alors en
|
||||||
|
attente.
|
||||||
|
|
||||||
:param client: Client RSS configuré, ou ``None`` lorsque le blog est désactivé.
|
:param client: Client RSS configuré, ou ``None`` lorsque le blog est désactivé.
|
||||||
:param state: État de déduplication et de cache HTTP associé au run.
|
:param state: État de déduplication et de cache HTTP associé au run.
|
||||||
:return: Nouveaux articles du blog.
|
:return: Résultat de récupération, incluant les métadonnées de cache.
|
||||||
:rtype: list[BlogArticle]
|
:rtype: BlogRSSFetchResult
|
||||||
:raises RuntimeError: Si la récupération RSS injectée échoue.
|
:raises RuntimeError: Si la récupération RSS injectée échoue.
|
||||||
"""
|
"""
|
||||||
if client is None or state is None:
|
if client is None or state is None:
|
||||||
return []
|
return BlogRSSFetchResult()
|
||||||
try:
|
try:
|
||||||
etag, last_modified = state.get_cache_headers()
|
etag, last_modified = state.get_cache_headers()
|
||||||
result = client.fetch_and_parse(
|
result = client.fetch_and_parse(
|
||||||
@@ -26,9 +31,8 @@ def fetch_blog_step(client: BlogRSSClient | None, state: BlogRSSState | None) ->
|
|||||||
)
|
)
|
||||||
if result.error is not None:
|
if result.error is not None:
|
||||||
raise RuntimeError(result.error) from None
|
raise RuntimeError(result.error) from None
|
||||||
if not result.not_modified:
|
if not result.not_modified and not result.articles:
|
||||||
state.add_guids(article.id for article in result.articles)
|
|
||||||
state.update_cache_headers(result.etag, result.last_modified)
|
state.update_cache_headers(result.etag, result.last_modified)
|
||||||
return list(result.articles)
|
return result
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
raise RuntimeError(f"Récupération du blog échouée : {redact_exception(exc)}") from None
|
raise RuntimeError(f"Récupération du blog échouée : {redact_exception(exc)}") from None
|
||||||
|
|||||||
@@ -19,6 +19,7 @@ import logging
|
|||||||
from collections.abc import Iterable
|
from collections.abc import Iterable
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
|
||||||
|
from pronote_sync.sources.blog.result import BlogRSSFetchResult
|
||||||
from pronote_sync.utils.redaction import redact_exception, redact_secrets
|
from pronote_sync.utils.redaction import redact_exception, redact_secrets
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
@@ -156,6 +157,22 @@ class BlogRSSState:
|
|||||||
self._known_guids.update(new_guids)
|
self._known_guids.update(new_guids)
|
||||||
self._save()
|
self._save()
|
||||||
|
|
||||||
|
def acknowledge(self, result: BlogRSSFetchResult) -> None:
|
||||||
|
"""Acquitte une récupération RSS après sa livraison confirmée.
|
||||||
|
|
||||||
|
Les GUID et les en-têtes de cache sont enregistrés ensemble afin qu'un
|
||||||
|
article dont la livraison a échoué reste récupérable à l'exécution
|
||||||
|
suivante. Une réponse ``304 Not Modified`` n'a rien à acquitter.
|
||||||
|
|
||||||
|
:param result: Résultat RSS livré avec succès.
|
||||||
|
"""
|
||||||
|
if result.not_modified:
|
||||||
|
return
|
||||||
|
self._known_guids.update(article.id for article in result.articles)
|
||||||
|
self._etag = result.etag
|
||||||
|
self._last_modified = result.last_modified
|
||||||
|
self._save()
|
||||||
|
|
||||||
def get_cache_headers(self) -> tuple[str | None, str | None]:
|
def get_cache_headers(self) -> tuple[str | None, str | None]:
|
||||||
"""Renvoie les en-têtes de cache HTTP mémorisés.
|
"""Renvoie les en-têtes de cache HTTP mémorisés.
|
||||||
|
|
||||||
|
|||||||
@@ -1040,6 +1040,110 @@ def test_runner_blog_success_delivers_articles_into_xmpp_message_external_info(
|
|||||||
assert xmpp_message.external_info.blog_articles[0].title == "Test Article"
|
assert xmpp_message.external_info.blog_articles[0].title == "Test Article"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize(
|
||||||
|
("channel_kind", "dry_run", "should_acknowledge"),
|
||||||
|
[
|
||||||
|
("success", False, True),
|
||||||
|
("false", False, False),
|
||||||
|
("exception", False, False),
|
||||||
|
("none", False, False),
|
||||||
|
("success", True, False),
|
||||||
|
],
|
||||||
|
)
|
||||||
|
def test_runner_acknowledges_blog_only_after_confirmed_xmpp_delivery(
|
||||||
|
pipeline_inputs: tuple[Lesson, Homework],
|
||||||
|
tmp_path: Any,
|
||||||
|
channel_kind: str,
|
||||||
|
dry_run: bool,
|
||||||
|
should_acknowledge: bool,
|
||||||
|
) -> None:
|
||||||
|
"""Les GUID RSS restent rejouables tant que XMPP n'a pas confirmé l'envoi.
|
||||||
|
|
||||||
|
:param pipeline_inputs: Données Pronote de test.
|
||||||
|
:param tmp_path: Répertoire temporaire pour l'état RSS.
|
||||||
|
:param channel_kind: Comportement du canal XMPP simulé.
|
||||||
|
:param dry_run: Active ou non le mode simulation.
|
||||||
|
:param should_acknowledge: Indique si l'état RSS doit être acquitté.
|
||||||
|
"""
|
||||||
|
lesson, homework = pipeline_inputs
|
||||||
|
calls: list[str] = []
|
||||||
|
state_file = tmp_path / "blog-state.json"
|
||||||
|
|
||||||
|
class SuccessfulBlogClient:
|
||||||
|
"""Client RSS renvoyant un article non encore livré."""
|
||||||
|
|
||||||
|
def fetch_and_parse(
|
||||||
|
self,
|
||||||
|
*,
|
||||||
|
known_guids: frozenset[str] | None = None,
|
||||||
|
etag: str | None = None,
|
||||||
|
last_modified: str | None = None,
|
||||||
|
) -> BlogRSSFetchResult:
|
||||||
|
"""Retourne un article et des en-têtes de cache déterministes.
|
||||||
|
|
||||||
|
:param known_guids: GUID déjà connus, ignorés dans ce faux client.
|
||||||
|
:param etag: ETag mémorisé, ignoré dans ce faux client.
|
||||||
|
:param last_modified: Date HTTP mémorisée, ignorée dans ce faux client.
|
||||||
|
:return: Résultat RSS avec un article à livrer.
|
||||||
|
:rtype: BlogRSSFetchResult
|
||||||
|
"""
|
||||||
|
del known_guids, etag, last_modified
|
||||||
|
return BlogRSSFetchResult(
|
||||||
|
articles=(
|
||||||
|
BlogArticle(
|
||||||
|
id="article-to-deliver",
|
||||||
|
title="Article à livrer",
|
||||||
|
url="https://example.com/article-to-deliver",
|
||||||
|
published_at=datetime(2026, 9, 8, 12, 0),
|
||||||
|
updated_at=None,
|
||||||
|
category=None,
|
||||||
|
author=None,
|
||||||
|
content_html="<p>Contenu</p>",
|
||||||
|
content_text="Contenu",
|
||||||
|
),
|
||||||
|
),
|
||||||
|
etag="etag-after-delivery",
|
||||||
|
last_modified="Tue, 08 Sep 2026 12:00:00 GMT",
|
||||||
|
)
|
||||||
|
|
||||||
|
channel: Any
|
||||||
|
if channel_kind == "success":
|
||||||
|
channel = StubChannel(calls)
|
||||||
|
elif channel_kind == "false":
|
||||||
|
channel = FailingChannel()
|
||||||
|
elif channel_kind == "exception":
|
||||||
|
channel = ExceptionalChannel()
|
||||||
|
else:
|
||||||
|
channel = None
|
||||||
|
|
||||||
|
runner = PipelineRunner(
|
||||||
|
settings=Settings(blog=Settings().blog.model_copy(update={"enabled": True})),
|
||||||
|
pronote_fetcher=StubFetcher(calls, lesson, homework),
|
||||||
|
caldav_synchronizer=lambda data, settings: successful_sync_result(),
|
||||||
|
agenda_comparator=cast("AgendaComparator | None", StubComparator(calls)),
|
||||||
|
blog_client=cast("BlogRSSClient | None", SuccessfulBlogClient()),
|
||||||
|
blog_state=BlogRSSState(state_file),
|
||||||
|
channel=channel,
|
||||||
|
dry_run=dry_run,
|
||||||
|
now_provider=lambda: datetime(2026, 9, 8, 7, 0),
|
||||||
|
)
|
||||||
|
|
||||||
|
data, errors = runner.run()
|
||||||
|
|
||||||
|
assert data is not None
|
||||||
|
if should_acknowledge:
|
||||||
|
acknowledged_state = BlogRSSState(state_file)
|
||||||
|
assert acknowledged_state.get_known_guids() == frozenset({"article-to-deliver"})
|
||||||
|
assert acknowledged_state.get_cache_headers() == (
|
||||||
|
"etag-after-delivery",
|
||||||
|
"Tue, 08 Sep 2026 12:00:00 GMT",
|
||||||
|
)
|
||||||
|
else:
|
||||||
|
assert not state_file.exists()
|
||||||
|
if channel_kind in {"false", "exception"}:
|
||||||
|
assert any(error.step == "send" for error in errors)
|
||||||
|
|
||||||
|
|
||||||
def test_runner_secret_redaction_in_pipeline_errors(
|
def test_runner_secret_redaction_in_pipeline_errors(
|
||||||
pipeline_inputs: tuple[Lesson, Homework],
|
pipeline_inputs: tuple[Lesson, Homework],
|
||||||
) -> None:
|
) -> None:
|
||||||
|
|||||||
@@ -14,11 +14,14 @@ Tous les tests utilisent des fichiers temporaires via la fixture ``tmp_path``.
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
|
from datetime import UTC, datetime
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
from unittest.mock import patch
|
from unittest.mock import patch
|
||||||
|
|
||||||
import pytest
|
import pytest
|
||||||
|
|
||||||
|
from pronote_sync.models.blog import BlogArticle
|
||||||
|
from pronote_sync.sources.blog.result import BlogRSSFetchResult
|
||||||
from pronote_sync.sources.blog.state import BlogRSSState
|
from pronote_sync.sources.blog.state import BlogRSSState
|
||||||
|
|
||||||
|
|
||||||
@@ -116,6 +119,38 @@ def test_add_guids_empty_noop(tmp_path: Path) -> None:
|
|||||||
assert state_file.read_text(encoding="utf-8") == original_content
|
assert state_file.read_text(encoding="utf-8") == original_content
|
||||||
|
|
||||||
|
|
||||||
|
def test_acknowledge_persists_guids_and_cache_headers_together(tmp_path: Path) -> None:
|
||||||
|
"""Vérifie l'acquittement atomique après une livraison confirmée.
|
||||||
|
|
||||||
|
:param tmp_path: Fixture pytest pour un répertoire temporaire.
|
||||||
|
:return: None
|
||||||
|
"""
|
||||||
|
state_file = tmp_path / "state.json"
|
||||||
|
state = BlogRSSState(state_file)
|
||||||
|
article = BlogArticle(
|
||||||
|
id="guid-1",
|
||||||
|
title="Article",
|
||||||
|
url="https://example.com/article",
|
||||||
|
published_at=datetime(2026, 9, 12, 8, 0, tzinfo=UTC),
|
||||||
|
updated_at=None,
|
||||||
|
category=None,
|
||||||
|
author=None,
|
||||||
|
content_html="<p>Contenu</p>",
|
||||||
|
content_text="Contenu",
|
||||||
|
)
|
||||||
|
|
||||||
|
state.acknowledge(
|
||||||
|
BlogRSSFetchResult(
|
||||||
|
articles=(article,),
|
||||||
|
etag="etag-1",
|
||||||
|
last_modified="Sat, 12 Sep 2026 08:00:00 GMT",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
assert state.get_known_guids() == frozenset({"guid-1"})
|
||||||
|
assert state.get_cache_headers() == ("etag-1", "Sat, 12 Sep 2026 08:00:00 GMT")
|
||||||
|
|
||||||
|
|
||||||
def test_state_load_persisted_guids(tmp_path: Path) -> None:
|
def test_state_load_persisted_guids(tmp_path: Path) -> None:
|
||||||
"""Vérifie que les GUID persistés sont rechargés dans une nouvelle instance.
|
"""Vérifie que les GUID persistés sont rechargés dans une nouvelle instance.
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user