Skip to content

Commit 4e3ff68

Browse files
authored
Merge pull request #45 from teams-notifier/feat/recover-missed-mr-close
feat: recover stalled MRs via emoji-triggered API refresh
2 parents fab4eaa + 82dc16b commit 4e3ff68

4 files changed

Lines changed: 749 additions & 1 deletion

File tree

db.py

Lines changed: 103 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -410,7 +410,7 @@ async def _generic_norm_upsert(
410410
INSERT INTO "{table}" (
411411
{", ".join(ins_col)}
412412
) VALUES (
413-
{", ".join(["$"+str(i+1) for i in range(len(ins_col))])}
413+
{", ".join(["$" + str(i + 1) for i in range(len(ins_col))])}
414414
) RETURNING {", ".join(sel_cols)}
415415
""",
416416
*ins_args,
@@ -495,6 +495,108 @@ async def get_pending_refreshes(self, limit: int = 50) -> list[dict[str, Any]]:
495495
)
496496
return [dict(row) for row in rows]
497497

498+
async def refresh_mr_payload_from_api(
499+
self, merge_request_ref_id: int, api_data: dict[str, Any]
500+
) -> MergeRequestInfos:
501+
"""Update stored MR payload with fresh state from GitLab API.
502+
503+
Syncs state, title, draft, merge status, branches, and pipeline ID.
504+
Assignees/reviewers are not updated (API response lacks email field required by GLUser).
505+
"""
506+
connection: asyncpg.Connection
507+
async with await database.acquire() as connection:
508+
async with connection.transaction():
509+
row = await connection.fetchrow(
510+
"""SELECT merge_request_ref_id, merge_request_payload,
511+
merge_request_extra_state, head_pipeline_id
512+
FROM merge_request_ref
513+
WHERE merge_request_ref_id = $1
514+
FOR UPDATE""",
515+
merge_request_ref_id,
516+
)
517+
assert row is not None
518+
519+
payload = row["merge_request_payload"]
520+
oa = payload.get("object_attributes", {})
521+
522+
for field in (
523+
"state",
524+
"title",
525+
"draft",
526+
"detailed_merge_status",
527+
"source_branch",
528+
"target_branch",
529+
):
530+
if field in api_data:
531+
oa[field] = api_data[field]
532+
533+
# GitLab REST API returns updated_at as ISO 8601 ("...Z");
534+
# webhook payloads use "YYYY-MM-DD HH:MM:SS UTC". Normalize
535+
# to webhook format so downstream fromisoformat parsing
536+
# (with the " UTC" -> "+00:00" replace) keeps working.
537+
# Defensive: if GitLab ever returns a naive datetime (no
538+
# offset), assume UTC rather than letting astimezone() apply
539+
# the host's local TZ. If parsing fails outright, log and
540+
# keep the stored value — never raise here, otherwise the
541+
# whole pending_mr_refresh row gets stuck retrying forever.
542+
if "updated_at" in api_data and api_data["updated_at"]:
543+
raw = api_data["updated_at"]
544+
try:
545+
parsed = datetime.datetime.fromisoformat(raw.replace("Z", "+00:00"))
546+
if parsed.tzinfo is None:
547+
parsed = parsed.replace(tzinfo=datetime.UTC)
548+
oa["updated_at"] = parsed.astimezone(datetime.UTC).strftime("%Y-%m-%d %H:%M:%S UTC")
549+
except (ValueError, TypeError) as exc:
550+
log.warning(
551+
"could not parse api updated_at, keeping stored value",
552+
merge_request_ref_id=merge_request_ref_id,
553+
raw=raw,
554+
error=str(exc),
555+
)
556+
557+
if "draft" in api_data:
558+
oa["work_in_progress"] = api_data["draft"]
559+
560+
# Synthesize `action` from state so cards/render.py picks the
561+
# right icon (CodeTextOff for close, Merge for merge). The
562+
# renderer keys off action, not state, so we MUST set it.
563+
# Only mutate on terminal/reopen transitions; otherwise keep
564+
# the webhook-recorded action to avoid masking real events.
565+
api_state = api_data.get("state")
566+
if api_state == "merged" and oa.get("action") != "merge":
567+
oa["action"] = "merge"
568+
elif api_state == "closed" and oa.get("action") != "close":
569+
oa["action"] = "close"
570+
elif api_state == "opened" and oa.get("action") in ("close", "merge"):
571+
oa["action"] = "reopen"
572+
573+
head_pipeline_id = row["head_pipeline_id"]
574+
api_pipeline = api_data.get("head_pipeline")
575+
if api_pipeline and api_pipeline.get("id"):
576+
oa["head_pipeline_id"] = api_pipeline["id"]
577+
head_pipeline_id = api_pipeline["id"]
578+
579+
payload["object_attributes"] = oa
580+
581+
await connection.execute(
582+
"""UPDATE merge_request_ref
583+
SET merge_request_payload = $1, head_pipeline_id = $2
584+
WHERE merge_request_ref_id = $3""",
585+
payload,
586+
head_pipeline_id,
587+
merge_request_ref_id,
588+
)
589+
590+
# extra_state is read from the pre-update `row` snapshot. Safe today
591+
# because this function does not mutate extra_state; if that ever
592+
# changes, re-read it after the UPDATE or RETURNING it.
593+
return MergeRequestInfos(
594+
merge_request_ref_id=merge_request_ref_id,
595+
merge_request_payload=payload,
596+
merge_request_extra_state=row["merge_request_extra_state"],
597+
head_pipeline_id=head_pipeline_id,
598+
)
599+
498600
async def delete_pending_refresh(self, merge_request_ref_id: int) -> None:
499601
"""Delete a pending refresh after processing."""
500602
connection: asyncpg.Connection

gitlab_api.py

Lines changed: 79 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515

1616
if TYPE_CHECKING:
1717
from db import MergeRequestExtraState
18+
from db import MergeRequestInfos
1819

1920
logger = fastapi_structured_logging.get_logger()
2021

@@ -178,3 +179,81 @@ async def _fetch_mr_discussion_stats(
178179
error=str(e),
179180
)
180181
return None
182+
183+
184+
async def fetch_and_refresh_mr_status(
185+
merge_request_ref_id: int,
186+
project_url: str,
187+
project_id: int,
188+
mr_iid: int,
189+
) -> MergeRequestInfos | None:
190+
"""
191+
Fetch current MR status from GitLab API and update stored payload.
192+
193+
Syncs state, title, draft, merge status, branches, and pipeline ID.
194+
Returns updated MergeRequestInfos or None if API unavailable
195+
(token missing, HTTP failure, or transient error).
196+
Callers should treat None as "could not verify" — distinct from
197+
"API confirmed MR still open".
198+
"""
199+
api_data = await _fetch_mr_status(project_url, project_id, mr_iid)
200+
if api_data is None:
201+
return None
202+
203+
from db import dbh
204+
205+
return await dbh.refresh_mr_payload_from_api(merge_request_ref_id, api_data)
206+
207+
208+
async def _fetch_mr_status(
209+
project_url: str,
210+
project_id: int,
211+
mr_iid: int,
212+
) -> dict[str, Any] | None:
213+
"""Fetch single MR from GitLab API. Returns None if token not configured or error."""
214+
api_token = config.get_gitlab_api_token(project_url)
215+
if api_token is None:
216+
logger.debug("no gitlab api token configured for project", project_url=project_url)
217+
return None
218+
219+
encoded_project_id = quote(str(project_id), safe="")
220+
url = f"{api_token.url.rstrip('/')}/api/v4/projects/{encoded_project_id}/merge_requests/{mr_iid}"
221+
222+
try:
223+
timeout = httpx.Timeout(5.0, connect=2.0)
224+
async with httpx.AsyncClient(timeout=timeout) as client:
225+
response = await client.get(
226+
url,
227+
headers={"PRIVATE-TOKEN": api_token.token},
228+
)
229+
response.raise_for_status()
230+
data: dict[str, Any] = response.json()
231+
logger.info(
232+
"fetched mr status from api",
233+
project_id=project_id,
234+
mr_iid=mr_iid,
235+
state=data.get("state"),
236+
draft=data.get("draft"),
237+
detailed_merge_status=data.get("detailed_merge_status"),
238+
)
239+
return data
240+
except httpx.HTTPStatusError as e:
241+
logger.warning(
242+
"gitlab api http error fetching mr status",
243+
token_name=api_token.name,
244+
api_url=url,
245+
project_id=project_id,
246+
mr_iid=mr_iid,
247+
status_code=e.response.status_code,
248+
)
249+
return None
250+
except Exception as e:
251+
logger.warning(
252+
"gitlab api error fetching mr status",
253+
token_name=api_token.name,
254+
api_url=url,
255+
project_id=project_id,
256+
mr_iid=mr_iid,
257+
error=str(e),
258+
)
259+
return None

periodic_cleanup.py

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,13 +8,15 @@
88

99
from cards.render import render
1010
from config import DefaultConfig
11+
from config import config
1112
from db import DatabaseLifecycleHandler
1213
from db import MergeRequestInfos
1314
from db import compute_mri_fingerprint
1415
from db import dbh
1516
from db import has_unresolved_threads
1617
from db import make_mr_summary
1718
from gitlab_api import fetch_and_persist_discussion_stats
19+
from gitlab_api import fetch_and_refresh_mr_status
1820
from webhook.messaging import update_all_messages_transactional
1921

2022

@@ -46,6 +48,41 @@ async def _process_pending_refreshes() -> int:
4648
head_pipeline_id=row["head_pipeline_id"],
4749
)
4850

51+
# Emoji events trigger a full MR status refresh via API
52+
# to sync state that may have been lost due to missed webhooks.
53+
# We do NOT short-circuit when stored state is already terminal:
54+
# the merge_request handler updates payload state and schedules
55+
# deletion in separate transactions, so a partial rollback can
56+
# leave state="merged" with refs still present. Letting the API
57+
# refresh + api-refresh-close branch run is the recovery for that.
58+
api_refreshed = False
59+
if row["payload_type"] == "emoji":
60+
refreshed_mri = await fetch_and_refresh_mr_status(
61+
merge_request_ref_id=mri.merge_request_ref_id,
62+
project_url=mri.merge_request_payload.project.web_url,
63+
project_id=mri.merge_request_payload.object_attributes.target_project_id,
64+
mr_iid=mri.merge_request_payload.object_attributes.iid,
65+
)
66+
if refreshed_mri is not None:
67+
mri = refreshed_mri
68+
api_refreshed = True
69+
if mri.merge_request_payload.object_attributes.state in ("closed", "merged"):
70+
logger.info(
71+
"api refresh detected terminal state",
72+
merge_request_ref_id=mri.merge_request_ref_id,
73+
state=mri.merge_request_payload.object_attributes.state,
74+
)
75+
else:
76+
# None = couldn't verify (no token / HTTP error). Distinct
77+
# from "API confirmed still open"; recurring emoji events
78+
# on a closed MR will keep landing here until API succeeds.
79+
logger.warning(
80+
"api refresh unavailable for emoji refresh",
81+
merge_request_ref_id=mri.merge_request_ref_id,
82+
project_id=mri.merge_request_payload.object_attributes.target_project_id,
83+
mr_iid=mri.merge_request_payload.object_attributes.iid,
84+
)
85+
4986
had_unresolved_threads = has_unresolved_threads(mri.merge_request_extra_state)
5087

5188
updated_extra_state = await fetch_and_persist_discussion_stats(
@@ -90,6 +127,29 @@ async def _process_pending_refreshes() -> int:
90127
)
91128
continue
92129

130+
# API refresh detected closed/merged MR — schedule message deletion to clean up
131+
if api_refreshed and is_closing_state:
132+
datasource_fingerprint = compute_mri_fingerprint(mri)
133+
card = render(mri, collapsed=True, show_collapsible=True)
134+
await update_all_messages_transactional(
135+
mri,
136+
card,
137+
make_mr_summary(mri),
138+
datasource_fingerprint,
139+
payload_updated_at,
140+
"api-refresh-close",
141+
schedule_deletion=True,
142+
deletion_delay=datetime.timedelta(seconds=config.MESSAGE_DELETE_DELAY_SECONDS),
143+
)
144+
await dbh.delete_pending_refresh(mri.merge_request_ref_id)
145+
processed += 1
146+
logger.info(
147+
"api refresh detected closed/merged MR, scheduled message deletion",
148+
merge_request_ref_id=mri.merge_request_ref_id,
149+
state=mri.merge_request_payload.object_attributes.state,
150+
)
151+
continue
152+
93153
should_be_collapsed: bool = (
94154
mri.merge_request_payload.object_attributes.draft
95155
or mri.merge_request_payload.object_attributes.work_in_progress

0 commit comments

Comments
 (0)