Skip to content

Commit 7fbfae4

Browse files
committed
fix(lark): reconcile a stall notice the provider already took
The transport records the provider locator of a sent message before it reads that message back, so a notice the provider accepted can come back `sent_unverified` with the text already on the channel. Counting that as a rejected notice let every later stalled attempt post the reader another copy, which is the one thing the notice is meant not to do. Keep the notice locator in the delivery record and verify it before sending again, exactly as a part is reconciled before it is re-sent, so the reader sees one notice per stalled sequence. The shared "the provider may already hold this text" predicate is the same one the part path uses, and the stall counter's comment now names what it counts. Signed-off-by: huangruiteng <14976749+huangruiteng@users.noreply.github.com>
1 parent 3105c56 commit 7fbfae4

2 files changed

Lines changed: 240 additions & 34 deletions

File tree

‎loopx/extensions/lark/manager_reply_parts.py‎

Lines changed: 149 additions & 34 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,11 @@
1515
A send the provider accepted but did not read back leaves its provider locator
1616
in the same record. The next attempt verifies that locator before sending the
1717
part again, so an ambiguous send is reconciled instead of repeated.
18+
19+
A partly delivered answer also has to look partly delivered: once the sequence
20+
has stopped mid-answer three times, one bounded notice tells the reader how much
21+
went out and where the rest is. That notice records and verifies its own
22+
locator the same way, so a reader is never told the same thing twice.
1823
"""
1924

2025
from __future__ import annotations
@@ -39,11 +44,13 @@
3944
PART_ATTEMPT_KEY = "delivery_part_attempt"
4045
PART_STALL_NOTICE_KEY = "delivery_part_stall_notice"
4146
PART_STALL_COUNT_KEY = "delivery_part_stall_count"
47+
PART_STALL_NOTICE_ATTEMPT_KEY = "delivery_part_stall_notice_attempt"
4248
PART_DELIVERY_INCOMPLETE = "reply_part_delivery_incomplete"
4349
PART_DELIVERY_COMPLETION_UNVERIFIED = "reply_part_delivery_completion_unverified"
44-
# The notice is a last resort, not the outcome: the sequence counts how often it
45-
# stopped mid-answer and only speaks after that happened more than once, so a
46-
# transient provider hiccup never turns into a message the reader did not need.
50+
# The notice is a last resort, not the outcome: the sequence counts the attempts
51+
# that stopped the answer mid-sequence (progress does not reset that count, only
52+
# a different split does) and only speaks after three of them, so a transient
53+
# provider hiccup never turns into a message the reader did not need.
4754
PART_STALL_NOTICE_MIN_STALLS = 3
4855
MANAGER_REPLY_STALL_NOTICE = (
4956
"本条答复超过可发送长度,目前只发出了前面的 {sent}/{count} 段;"
@@ -66,13 +73,13 @@ def plan_manager_reply_parts(reply_text: str) -> tuple[list[str], bool]:
6673
return parts, truncated
6774

6875

69-
def _part_verified(reply: Mapping[str, Any]) -> bool:
70-
"""Whether the provider reported this part present on the channel.
76+
def _reply_verified(reply: Mapping[str, Any]) -> bool:
77+
"""Whether the provider reported this reply present on the channel.
7178
72-
A part can be confirmed either by the readback that follows its own send or
79+
A reply can be confirmed either by the readback that follows its own send or
7380
by the reconciliation of an earlier send the provider accepted but did not
74-
read back. Both mean the reader has that text, which is the fact the counter
75-
records; a reconciled part has no new write of its own.
81+
read back. Both mean the reader has that text, which is the fact the record
82+
keeps; a reconciled send has no new write of its own.
7683
"""
7784

7885
return bool(
@@ -81,25 +88,31 @@ def _part_verified(reply: Mapping[str, Any]) -> bool:
8188
)
8289

8390

84-
def _part_accepted(reply: Mapping[str, Any]) -> bool:
85-
"""Whether this part may be counted as delivered.
91+
def _reply_on_channel(reply: Mapping[str, Any]) -> bool:
92+
"""Whether the reader may already have this reply's text.
8693
87-
``ok`` also requires the source reaction cleanup to have finished, so a part
88-
the provider already verified can come back not-ok with a cleanup still
89-
pending. Its text is on the channel either way: counting it is what keeps a
90-
retry from sending the reader the same part twice, and the pending cleanup
91-
stays the transport's own business.
94+
``ok`` also requires the source reaction cleanup to have finished, so a
95+
reply the provider already accepted can come back not-ok with a cleanup
96+
still pending. Treating that as delivered is what keeps a retry from sending
97+
the reader the same text twice, and the pending cleanup stays the
98+
transport's own business.
9299
"""
93100

94-
return reply.get("ok") is True or _part_verified(reply)
101+
return reply.get("ok") is True or _reply_verified(reply)
102+
103+
104+
def _part_accepted(reply: Mapping[str, Any]) -> bool:
105+
"""Whether this part may be counted as delivered."""
106+
107+
return _reply_on_channel(reply)
95108

96109

97110
def _accepted_reply_facts(reply: Mapping[str, Any]) -> dict[str, Any]:
98111
"""The durable facts of one part the provider confirmed."""
99112

100113
return {
101114
"reply_idempotency_key": reply.get("idempotency_key"),
102-
PART_DELIVERY_VERIFIED_KEY: _part_verified(reply),
115+
PART_DELIVERY_VERIFIED_KEY: _reply_verified(reply),
103116
}
104117

105118

@@ -156,6 +169,15 @@ def part_delivery_incomplete_reason(delivery_state: Mapping[str, Any]) -> str:
156169
return PART_DELIVERY_INCOMPLETE
157170

158171

172+
def _recorded_attempt(recorded: Any) -> Mapping[str, Any] | None:
173+
"""The provider locator inside one recorded attempt, when it is well formed."""
174+
175+
if not isinstance(recorded, Mapping):
176+
return None
177+
attempt = recorded.get("attempt")
178+
return attempt if isinstance(attempt, Mapping) else None
179+
180+
159181
def recorded_part_attempt(
160182
delivery_state: Mapping[str, Any], index: int
161183
) -> Mapping[str, Any] | None:
@@ -164,8 +186,15 @@ def recorded_part_attempt(
164186
recorded = delivery_state.get(PART_ATTEMPT_KEY)
165187
if not isinstance(recorded, Mapping) or recorded.get("index") != index:
166188
return None
167-
attempt = recorded.get("attempt")
168-
return attempt if isinstance(attempt, Mapping) else None
189+
return _recorded_attempt(recorded)
190+
191+
192+
def recorded_stall_notice_attempt(
193+
delivery_state: Mapping[str, Any],
194+
) -> Mapping[str, Any] | None:
195+
"""The provider locator of a stall notice that was sent but not confirmed."""
196+
197+
return _recorded_attempt(delivery_state.get(PART_STALL_NOTICE_ATTEMPT_KEY))
169198

170199

171200
def reconciled_part_reply(
@@ -202,14 +231,47 @@ def reconciled_part_reply(
202231
return {**dict(verified), "part_reconciled": True}
203232

204233

234+
def reconciled_stall_notice(
235+
*,
236+
notice: str,
237+
delivery_state: Mapping[str, Any],
238+
reply_runner: Any,
239+
root: Path,
240+
config_path: Path,
241+
message_id: str,
242+
) -> Mapping[str, Any] | None:
243+
"""Confirm a previously sent stall notice instead of posting it twice.
244+
245+
The provider can accept the notice and still fail the readback that proves
246+
it, and this sequence keeps retrying until the remaining parts go through.
247+
The recorded locator is what keeps such a notice from being sent again on
248+
every later attempt, exactly as a part is reconciled before it is re-sent.
249+
"""
250+
251+
attempt = recorded_stall_notice_attempt(delivery_state)
252+
if attempt is None:
253+
return None
254+
verified = verify_lark_inbox_reply(
255+
project=root,
256+
config_path=config_path,
257+
message_id=message_id,
258+
text=notice,
259+
attempt=attempt,
260+
runner=reply_runner,
261+
)
262+
if verified.get("reply_verified") is not True:
263+
return None
264+
return {**dict(verified), "notice_reconciled": True}
265+
266+
205267
def plan_stalled_part_notice(delivery_state: Mapping[str, Any]) -> str | None:
206268
"""The bounded notice a stalled sequence posts once, after real retries.
207269
208270
A reader who received the first parts of an over-limit answer currently
209271
learns nothing more: the overflow note that says where the full answer lives
210272
only travels with the last part, and the remaining parts are retried in the
211273
background. This notice states what was delivered and where the rest is, and
212-
it waits for more than one failed attempt so an ordinary hiccup stays quiet.
274+
it waits for three stalled attempts so an ordinary hiccup stays quiet.
213275
"""
214276

215277
if delivery_state.get(PART_STALL_NOTICE_KEY) is True:
@@ -231,6 +293,58 @@ def plan_stalled_part_notice(delivery_state: Mapping[str, Any]) -> str | None:
231293
return MANAGER_REPLY_STALL_NOTICE.format(sent=sent, count=count)
232294

233295

296+
def deliver_stall_notice(
297+
*,
298+
delivery_state: dict[str, Any],
299+
delivery_path: Path,
300+
write_delivery,
301+
reply_runner: Any,
302+
root: Path,
303+
config_path: Path,
304+
message_id: str,
305+
) -> Mapping[str, Any] | None:
306+
"""Post the once-per-sequence stall notice, confirming a prior send first.
307+
308+
Returns the transport result of the send or reconciliation, and ``None``
309+
when this sequence has nothing to tell the reader.
310+
"""
311+
312+
notice = plan_stalled_part_notice(delivery_state)
313+
if notice is None:
314+
return None
315+
reconciled = reconciled_stall_notice(
316+
notice=notice,
317+
delivery_state=delivery_state,
318+
reply_runner=reply_runner,
319+
root=root,
320+
config_path=config_path,
321+
message_id=message_id,
322+
)
323+
# The locator of the send being attempted now replaces any older one, so the
324+
# record always points at the most recent unconfirmed notice.
325+
delivery_state.pop(PART_STALL_NOTICE_ATTEMPT_KEY, None)
326+
if reconciled is not None:
327+
return reconciled
328+
329+
def record_attempt(attempt: Mapping[str, Any]) -> None:
330+
# The locator has to survive the attempt that produced it: a retry
331+
# reloads the record and verifies it instead of posting the notice again.
332+
delivery_state[PART_STALL_NOTICE_ATTEMPT_KEY] = {"attempt": dict(attempt)}
333+
delivery_state["updated_at"] = datetime.now(timezone.utc).isoformat()
334+
write_delivery(delivery_path, delivery_state)
335+
336+
return reply_lark_event_inbox(
337+
project=root,
338+
config_path=config_path,
339+
message_id=message_id,
340+
text=notice,
341+
content_format="text",
342+
execute=True,
343+
runner=reply_runner,
344+
delivery_attempt_recorder=record_attempt,
345+
)
346+
347+
234348
def deliver_manager_reply_parts(
235349
*,
236350
parts: list[str],
@@ -355,8 +469,10 @@ def deliver_manager_reply_after_length_failure(
355469
"""Deliver one over-limit manager answer as bounded parts.
356470
357471
Returns the last accepted reply, or ``None`` plus the reason to report when
358-
a part was rejected. Plain text is the only format a split can promise, so
359-
the caller has already degraded presentation before calling this.
472+
a part was rejected. A sequence that stops with parts still unsent also tells
473+
the reader what went out, once, after enough failed attempts. Plain text is
474+
the only format a split can promise, so the caller has already degraded
475+
presentation before calling this.
360476
"""
361477

362478
parts, truncated = plan_manager_reply_parts(reply_text)
@@ -384,18 +500,17 @@ def deliver_manager_reply_after_length_failure(
384500
delivery_state[PART_STALL_COUNT_KEY] = (
385501
int(delivery_state.get(PART_STALL_COUNT_KEY) or 0) + 1
386502
)
387-
notice = plan_stalled_part_notice(delivery_state)
388-
if notice is not None:
389-
spoken = reply_lark_event_inbox(
390-
project=root,
391-
config_path=config_path,
392-
message_id=message_id,
393-
text=notice,
394-
content_format="text",
395-
execute=True,
396-
runner=reply_runner,
397-
)
398-
if spoken.get("ok") is True or spoken.get("reply_verified") is True:
503+
spoken = deliver_stall_notice(
504+
delivery_state=delivery_state,
505+
delivery_path=delivery_path,
506+
write_delivery=write_delivery,
507+
reply_runner=reply_runner,
508+
root=root,
509+
config_path=config_path,
510+
message_id=message_id,
511+
)
512+
if spoken is not None:
513+
if _reply_on_channel(spoken):
399514
delivery_state[PART_STALL_NOTICE_KEY] = True
400515
delivery_state["last_delivery_notice_status"] = str(
401516
spoken.get("status") or "reply_failed"

‎tests/extensions/test_lark_manager_reply_parts.py‎

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@
1515
PART_DELIVERY_COMPLETION_UNVERIFIED,
1616
PART_DELIVERY_INCOMPLETE,
1717
PART_DELIVERY_VERIFIED_KEY,
18+
PART_STALL_NOTICE_ATTEMPT_KEY,
1819
PART_STALL_COUNT_KEY,
1920
PART_STALL_NOTICE_KEY,
2021
completed_part_delivery_receipt,
@@ -24,6 +25,7 @@
2425
plan_manager_reply_parts,
2526
plan_stalled_part_notice,
2627
recorded_part_attempt,
28+
recorded_stall_notice_attempt,
2729
)
2830

2931

@@ -494,3 +496,92 @@ def rejecting_everything(**kwargs):
494496
message_id="om_fixture",
495497
)
496498
assert any(text.startswith("本条答复") for text in attempts)
499+
500+
501+
def test_a_notice_the_provider_took_but_did_not_read_back_is_confirmed(
502+
monkeypatch, delivery,
503+
):
504+
"""A notice the provider accepted but could not read back is not repeated.
505+
506+
The transport records the provider locator of the notice before it reads the
507+
message back, so a send that comes back ``sent_unverified`` has to be
508+
confirmed on the next attempt instead of being posted to the reader twice.
509+
"""
510+
511+
parts, _ = plan_manager_reply_parts(BODY)
512+
state = _stalled_state(parts, stalls=2, sent=2)
513+
notice = MANAGER_REPLY_STALL_NOTICE.format(sent=2, count=len(parts))
514+
sends: list[str] = []
515+
verifications: list[str] = []
516+
517+
def ambiguous_send(**kwargs):
518+
sends.append(kwargs["text"])
519+
if kwargs["text"].startswith("("):
520+
return {
521+
"ok": False,
522+
"status": "reply_provider_failed",
523+
"idempotency_key": None,
524+
}
525+
kwargs["delivery_attempt_recorder"](
526+
{
527+
"schema_version": "manager_return_delivery_attempt_v0",
528+
"provider": "lark",
529+
"message_ref": "om_notice",
530+
"intent_digest": "sha256:notice-intent",
531+
"provider_receipt": "sha256:notice-receipt",
532+
}
533+
)
534+
return {
535+
"ok": False,
536+
"status": "sent_unverified",
537+
"idempotency_key": "sha256:notice-receipt",
538+
"write_performed": True,
539+
"verification_performed": True,
540+
"reply_verified": False,
541+
}
542+
543+
def confirmed_notice(**kwargs):
544+
verifications.append(kwargs["text"])
545+
return {
546+
"ok": True,
547+
"status": "sent_verified",
548+
"idempotency_key": "sha256:notice-receipt",
549+
"verification_performed": True,
550+
"reply_verified": True,
551+
}
552+
553+
monkeypatch.setattr(parts_module, "reply_lark_event_inbox", ambiguous_send)
554+
monkeypatch.setattr(parts_module, "verify_lark_inbox_reply", confirmed_notice)
555+
556+
def deliver_with_stalls():
557+
return deliver_manager_reply_after_length_failure(
558+
reply_text=BODY,
559+
delivery_state=state,
560+
delivery_path=delivery["tmp"] / "delivery.json",
561+
write_delivery=lambda path, payload: None,
562+
reply_runner=object(),
563+
root=delivery["tmp"],
564+
config_path=delivery["tmp"] / "config.json",
565+
message_id="om_fixture",
566+
)
567+
568+
deliver_with_stalls()
569+
570+
# The notice went out but its readback did not confirm it, so the reader may
571+
# already have it and the record keeps that attempt's locator.
572+
assert [text for text in sends if text.startswith("本条答复")] == [notice]
573+
assert state.get(PART_STALL_NOTICE_KEY) is not True
574+
assert recorded_stall_notice_attempt(state) == {
575+
"schema_version": "manager_return_delivery_attempt_v0",
576+
"provider": "lark",
577+
"message_ref": "om_notice",
578+
"intent_digest": "sha256:notice-intent",
579+
"provider_receipt": "sha256:notice-receipt",
580+
}
581+
582+
deliver_with_stalls()
583+
584+
assert verifications == [notice]
585+
assert [text for text in sends if text.startswith("本条答复")] == [notice]
586+
assert state[PART_STALL_NOTICE_KEY] is True
587+
assert PART_STALL_NOTICE_ATTEMPT_KEY not in state

0 commit comments

Comments
 (0)