feat: improve LXMF message handling with new oracles and progress polling improvements

This commit is contained in:
Ivan 2026-07-24 08:57:17 -05:00
parent 58f036219a
commit 1212266be5
No known key found for this signature in database
9 changed files with 555 additions and 20 deletions

View file

@ -39,6 +39,7 @@ All notable changes to this project will be documented in this file.
- Bundled [RNS-over-HTTP](https://github.com/Quad4-Software/RNS-over-HTTP) HTTPInterface with Interfaces page client/server setup, auto-install into the Reticulum interface path, and httpx support (Android includes httpx with HTTP/2 as pure-Python wheels)
- Docker **extra** image variant (`-extra` suffix): same Alpine Dockerfile with `VARIANT=extra` (adds **i2pd** and **yggdrasil**), built and published beside standard and hardened
- Docker Alpine/hardened recipes moved into `scripts/docker/*.sh` so Dockerfiles stay thin (shared frontend, venv, overlay, and runtime setup)
- LXMF send/delivery oracles for no-proof SENT retries, auto-resend claim rules, and ConversationViewer single-bubble absorb/created races
### Changed
@ -68,6 +69,7 @@ All notable changes to this project will be documented in this file.
### Fixed
- LXMF outbound progress polling stops on REJECTED (previously only DELIVERED / propagated SENT / FAILED / CANCELLED stopped the loop)
- Web Sync Messages after a backgrounded browser tab: recover stale WebSocket as a shell reconnect, refresh CSRF, and do not abort sync when request-path priming fails
- Conversations: re-opening an already-read thread no longer decrements the Messages unread badge
- Notifications: DND still updates the Messages unread badge (DND only suppresses OS notifications and sound)

Binary file not shown.

View file

@ -151,6 +151,7 @@ from meshchatx.src.backend.lxmf_utils import (
convert_lxmf_message_to_dict,
convert_lxmf_method_to_string,
convert_lxmf_state_to_string,
is_lxmf_outbound_progress_terminal,
is_user_facing_lxmf_payload,
lxmf_fields_are_reaction,
lxmf_sidebar_preview_for_conversation_latest_row,
@ -9531,20 +9532,8 @@ class ReticulumMeshChat:
),
)
# check message state
has_delivered = lxmf_message.state == LXMF.LXMessage.DELIVERED
has_propagated = (
lxmf_message.state == LXMF.LXMessage.SENT
and lxmf_message.method == LXMF.LXMessage.PROPAGATED
)
has_failed = (
lxmf_message.state == LXMF.LXMessage.FAILED
and getattr(lxmf_message, "try_propagation_on_fail", False) is not True
)
is_cancelled = lxmf_message.state == LXMF.LXMessage.CANCELLED
# check if we should stop updating
if has_delivered or has_propagated or has_failed or is_cancelled:
if is_lxmf_outbound_progress_terminal(lxmf_message):
should_update_message = False
else:
await asyncio.sleep(1)

View file

@ -806,6 +806,39 @@ def convert_lxmf_message_to_dict(
return out
def is_lxmf_outbound_progress_terminal(lxmf_message: LXMF.LXMessage) -> bool:
"""True when MeshChat should stop polling outbound progress.
Opportunistic or direct SENT is not terminal. Without an RNS delivery proof
the message stays SENT while LXMF retries, and only later becomes
DELIVERED, FAILED, or CANCELLED. Propagated SENT means parked on a node.
FAILED with try_propagation_on_fail still pending is also non-terminal.
REJECTED is terminal (peer refused the message).
"""
if lxmf_message.state == LXMF.LXMessage.DELIVERED:
return True
if lxmf_message.state == LXMF.LXMessage.REJECTED:
return True
if (
lxmf_message.state == LXMF.LXMessage.SENT
and lxmf_message.method == LXMF.LXMessage.PROPAGATED
):
return True
if lxmf_message.state == LXMF.LXMessage.CANCELLED:
return True
if (
lxmf_message.state == LXMF.LXMessage.FAILED
and getattr(
lxmf_message,
"try_propagation_on_fail",
False,
)
is not True
):
return True
return False
def convert_lxmf_state_to_string(lxmf_message: LXMF.LXMessage):
# convert state to string
lxmf_message_state = "unknown"

View file

@ -0,0 +1,461 @@
# SPDX-License-Identifier: 0BSD
"""Oracles for LXMF outbound send, delivery proofs, and auto-resend cascades.
Guarantee: missing RNS delivery proofs leave opportunistic/direct messages in
SENT (not DELIVERED). Progress polling must not treat that as terminal. Auto
resend only claims failed rows and emits a new hash. Same-content guards do not
treat a failed-only lineage as recent success.
These tests document the no-receipt double-send cascade:
LXMF retries the same hash until MAX_DELIVERY_ATTEMPTS, then MeshChat may
auto-resend a new hash when a peer announces or a path/ping succeeds.
"""
from __future__ import annotations
import asyncio
import time
import types
from types import SimpleNamespace
from unittest.mock import AsyncMock, MagicMock, patch
import LXMF
import pytest
from hypothesis import given, settings
from hypothesis import strategies as st
from meshchatx.meshchat import ReticulumMeshChat
from meshchatx.src.backend import auto_resend_guard as guard
from meshchatx.src.backend.database import Database
from meshchatx.src.backend.database.provider import DatabaseProvider
from meshchatx.src.backend.database.schema import DatabaseSchema
from meshchatx.src.backend.lxmf_utils import is_lxmf_outbound_progress_terminal
def _msg(**attrs):
defaults = {
"state": LXMF.LXMessage.OUTBOUND,
"method": LXMF.LXMessage.DIRECT,
"try_propagation_on_fail": False,
}
defaults.update(attrs)
return SimpleNamespace(**defaults)
@pytest.fixture
def db(tmp_path):
path = str(tmp_path / "send_delivery_oracle.db")
provider = DatabaseProvider(path)
DatabaseSchema(provider).initialize()
database = Database(path)
yield database
database.close_all()
provider.close_all()
def _insert_row(
db,
*,
msg_hash,
peer,
content,
state,
fields="{}",
delivery_attempts=0,
ts=None,
next_at=None,
):
now = ts if ts is not None else time.time()
db.messages.upsert_lxmf_message(
{
"hash": msg_hash,
"source_hash": peer,
"destination_hash": peer,
"peer_hash": peer,
"state": state,
"progress": 1 if state == "delivered" else 0,
"is_incoming": 0,
"method": "opportunistic",
"delivery_attempts": delivery_attempts,
"next_delivery_attempt_at": next_at,
"title": "",
"content": content,
"fields": fields,
"rssi": None,
"snr": None,
"quality": None,
"is_spam": 0,
"reply_to_hash": None,
"attachments_stripped": 0,
"timestamp": now,
},
)
def _bind_resend_app(db):
app = MagicMock()
app._auto_resend_coordinator = guard.AutoResendCoordinator()
app.websocket_broadcast = AsyncMock()
app.send_message = AsyncMock()
app.resend_failed_messages_for_destination = types.MethodType(
ReticulumMeshChat.resend_failed_messages_for_destination,
app,
)
ctx = MagicMock()
ctx.identity.hash.hex.return_value = "aa" * 16
ctx.database = db
ctx.config.allow_auto_resending_failed_messages_with_attachments.get.return_value = True
return app, ctx
def test_oracle_lxmf_sent_is_not_delivered_without_proof():
"""Independent LXMF contract: SENT != DELIVERED. Proofs flip to DELIVERED."""
assert LXMF.LXMessage.SENT != LXMF.LXMessage.DELIVERED
assert LXMF.LXMRouter.MAX_DELIVERY_ATTEMPTS >= 1
assert LXMF.LXMRouter.DELIVERY_RETRY_WAIT > 0
@pytest.mark.parametrize(
("state", "method", "try_prop", "expected"),
[
(LXMF.LXMessage.SENT, LXMF.LXMessage.OPPORTUNISTIC, False, False),
(LXMF.LXMessage.SENT, LXMF.LXMessage.DIRECT, False, False),
(LXMF.LXMessage.OUTBOUND, LXMF.LXMessage.DIRECT, False, False),
(LXMF.LXMessage.SENDING, LXMF.LXMessage.DIRECT, False, False),
(LXMF.LXMessage.SENT, LXMF.LXMessage.PROPAGATED, False, True),
(LXMF.LXMessage.DELIVERED, LXMF.LXMessage.OPPORTUNISTIC, False, True),
(LXMF.LXMessage.DELIVERED, LXMF.LXMessage.PROPAGATED, False, True),
(LXMF.LXMessage.FAILED, LXMF.LXMessage.DIRECT, False, True),
(LXMF.LXMessage.FAILED, LXMF.LXMessage.DIRECT, True, False),
(LXMF.LXMessage.CANCELLED, LXMF.LXMessage.DIRECT, False, True),
(LXMF.LXMessage.REJECTED, LXMF.LXMessage.DIRECT, False, True),
],
)
def test_oracle_progress_terminal_matrix(state, method, try_prop, expected):
msg = _msg(state=state, method=method, try_propagation_on_fail=try_prop)
assert is_lxmf_outbound_progress_terminal(msg) is expected
@given(
method=st.sampled_from(
[
LXMF.LXMessage.OPPORTUNISTIC,
LXMF.LXMessage.DIRECT,
LXMF.LXMessage.PAPER,
]
)
)
@settings(max_examples=24, deadline=None)
def test_oracle_sent_without_propagated_never_terminal(method):
"""No-receipt peers leave opportunistic/direct at SENT. Must keep polling."""
assert not is_lxmf_outbound_progress_terminal(
_msg(state=LXMF.LXMessage.SENT, method=method)
)
@pytest.mark.asyncio
async def test_oracle_progress_loop_keeps_polling_opportunistic_sent(db):
app = ReticulumMeshChat.__new__(ReticulumMeshChat)
app.current_context = MagicMock()
app.current_context.database = db
app.reticulum = None
app.websocket_broadcast = AsyncMock()
app._maybe_store_path_at_send_for_lxmf = MagicMock()
app._merge_stored_path_fields_from_db = MagicMock()
peer = "b" * 32
msg_hash = "c" * 32
_insert_row(db, msg_hash=msg_hash, peer=peer, content="hi", state="outbound")
mock_msg = MagicMock()
mock_msg.hash = bytes.fromhex(msg_hash)
mock_msg.progress = 0.0
mock_msg.delivery_attempts = 1
mock_msg.state = LXMF.LXMessage.SENT
mock_msg.method = LXMF.LXMessage.OPPORTUNISTIC
mock_msg.try_propagation_on_fail = False
mock_msg.next_delivery_attempt = None
# Avoid MagicMock auto-attrs leaking into SQLite binds.
type(mock_msg).rssi = property(lambda self: None)
type(mock_msg).snr = property(lambda self: None)
type(mock_msg).q = property(lambda self: None)
sleeps = 0
async def count_sleep(_duration):
nonlocal sleeps
sleeps += 1
if sleeps >= 2:
mock_msg.state = LXMF.LXMessage.FAILED
with (
patch(
"meshchatx.meshchat.convert_lxmf_message_to_dict",
return_value={"hash": msg_hash, "state": "sent"},
),
patch("asyncio.sleep", side_effect=count_sleep),
):
await ReticulumMeshChat.handle_lxmf_message_progress(
app,
mock_msg,
context=app.current_context,
)
assert sleeps >= 2
row = db.provider.fetchone(
"SELECT state FROM lxmf_messages WHERE hash = ?",
(msg_hash,),
)
assert row["state"] == "failed"
def test_oracle_claim_refuses_non_failed_states(db):
peer = "d" * 32
now = time.time()
for state, msg_hash in (
("sent", "1" * 32),
("outbound", "2" * 32),
("delivered", "3" * 32),
("sending", "4" * 32),
):
_insert_row(
db,
msg_hash=msg_hash,
peer=peer,
content="body",
state=state,
delivery_attempts=3,
)
assert not db.messages.try_claim_failed_message_for_auto_resend(
msg_hash,
cooldown_until=now + 120,
now=now,
)
def test_oracle_failed_only_lineage_does_not_count_as_recent_success(db):
"""After no-proof FAILED, only the failed row remains. Auto-resend may fire."""
peer = "e" * 32
now = time.time()
_insert_row(
db,
msg_hash="5" * 32,
peer=peer,
content="same body",
state="failed",
delivery_attempts=LXMF.LXMRouter.MAX_DELIVERY_ATTEMPTS,
ts=now - 10,
)
assert not db.messages.has_recent_outbound_with_content(
peer,
"same body",
within_seconds=300,
now=now,
)
def test_oracle_in_flight_sent_blocks_same_content_resend(db):
"""While LXMF still holds SENT (retrying proofs), do not treat as free to resend."""
peer = "f" * 32
now = time.time()
_insert_row(
db,
msg_hash="6" * 32,
peer=peer,
content="in flight",
state="sent",
delivery_attempts=2,
ts=now - 5,
)
assert db.messages.has_recent_outbound_with_content(
peer,
"in flight",
within_seconds=300,
now=now,
)
@pytest.mark.asyncio
async def test_oracle_auto_resend_skips_sent_rows_and_only_claims_failed(db):
peer = "7" * 32
_insert_row(db, msg_hash="8" * 32, peer=peer, content="alive", state="sent")
_insert_row(db, msg_hash="9" * 32, peer=peer, content="dead", state="failed")
app, ctx = _bind_resend_app(db)
new_hash = bytes.fromhex("a" * 32)
async def send_ok(*args, **kwargs):
m = MagicMock()
m.hash = new_hash
_insert_row(
db,
msg_hash=new_hash.hex(),
peer=peer,
content="dead",
state="outbound",
)
return m
app.send_message.side_effect = send_ok
await app.resend_failed_messages_for_destination(peer, context=ctx)
assert app.send_message.await_count == 1
assert (
db.provider.fetchone(
"SELECT state FROM lxmf_messages WHERE hash = ?",
("8" * 32,),
)["state"]
== "sent"
)
assert (
db.provider.fetchone(
"SELECT 1 FROM lxmf_messages WHERE hash = ?",
("9" * 32,),
)
is None
)
@pytest.mark.asyncio
async def test_oracle_auto_resend_emits_new_hash_not_same_hash(db):
peer = "b" * 32
old_hash = "c" * 32
new_hash = bytes.fromhex("d" * 32)
_insert_row(db, msg_hash=old_hash, peer=peer, content="retry", state="failed")
app, ctx = _bind_resend_app(db)
async def send_ok(*args, **kwargs):
m = MagicMock()
m.hash = new_hash
_insert_row(
db,
msg_hash=new_hash.hex(),
peer=peer,
content="retry",
state="outbound",
)
return m
app.send_message.side_effect = send_ok
await app.resend_failed_messages_for_destination(peer, context=ctx)
assert new_hash.hex() != old_hash
assert app.send_message.await_count == 1
assert (
guard.read_auto_resend_count(
db.provider.fetchone(
"SELECT fields FROM lxmf_messages WHERE hash = ?",
(new_hash.hex(),),
)["fields"]
)
== 1
)
@pytest.mark.asyncio
async def test_oracle_propagation_fallback_clears_packed_and_requeues_same_object():
app = ReticulumMeshChat.__new__(ReticulumMeshChat)
ctx = MagicMock()
ctx.forwarding_manager = None
ctx.message_router = MagicMock()
app.current_context = ctx
original_id = object()
lxmf_message = MagicMock()
lxmf_message.__id = original_id
lxmf_message.packed = b"stale-packet"
lxmf_message.delivery_attempts = 5
lxmf_message.next_delivery_attempt = time.time() + 10
lxmf_message.source_hash = bytes.fromhex("11" * 16)
lxmf_message.message_id = b"mid" * 5 + b"xx"
ReticulumMeshChat.send_failed_message_via_propagation_node(
app,
lxmf_message,
context=ctx,
)
assert lxmf_message.packed is None
assert lxmf_message.delivery_attempts == 0
assert lxmf_message.desired_method == LXMF.LXMessage.PROPAGATED
assert lxmf_message.try_propagation_on_fail is False
assert not hasattr(lxmf_message, "next_delivery_attempt")
ctx.message_router.handle_outbound.assert_called_once_with(lxmf_message)
@pytest.mark.asyncio
async def test_oracle_concurrent_resend_claims_send_once(db):
peer = "e" * 32
old_hash = "f" * 32
_insert_row(db, msg_hash=old_hash, peer=peer, content="race", state="failed")
app, ctx = _bind_resend_app(db)
async def slow_send(*args, **kwargs):
await asyncio.sleep(0.05)
m = MagicMock()
m.hash = bytes.fromhex("1" * 32)
_insert_row(
db,
msg_hash=m.hash.hex(),
peer=peer,
content="race",
state="outbound",
)
return m
app.send_message.side_effect = slow_send
await asyncio.gather(
app.resend_failed_messages_for_destination(peer, context=ctx),
app.resend_failed_messages_for_destination(peer, context=ctx),
app.resend_failed_messages_for_destination(peer, context=ctx),
)
assert app.send_message.await_count == 1
def test_oracle_no_receipt_cascade_contract_documented():
"""Closed contract for peers that never return delivery proofs."""
from pathlib import Path
# 1) Proof required for DELIVERED
assert LXMF.LXMessage.SENT != LXMF.LXMessage.DELIVERED
# 2) Library retries same outbound object
assert LXMF.LXMRouter.MAX_DELIVERY_ATTEMPTS == 5
# 3) Progress must keep polling SENT opportunistic
assert not is_lxmf_outbound_progress_terminal(
_msg(
state=LXMF.LXMessage.SENT,
method=LXMF.LXMessage.OPPORTUNISTIC,
)
)
# 4) Only FAILED is auto-resend eligible (claim SQL requires state=failed)
messages_src = (
Path(__file__).resolve().parents[2]
/ "meshchatx"
/ "src"
/ "backend"
/ "database"
/ "messages.py"
).read_text(encoding="utf-8")
assert "state = 'failed'" in messages_src
assert "try_claim_failed_message_for_auto_resend" in messages_src
# 5) Auto-resend budget caps replacement floods
assert guard.MAX_AUTO_RESEND_ATTEMPTS == 3
assert guard.AUTO_RESEND_COOLDOWN_SECONDS == 120
assert guard.RECENT_SAME_CONTENT_SECONDS == 300
def test_lxmf_send_delivery_oracle_proved():
pairs = 0
for state, method, expected in (
(LXMF.LXMessage.SENT, LXMF.LXMessage.OPPORTUNISTIC, False),
(LXMF.LXMessage.SENT, LXMF.LXMessage.DIRECT, False),
(LXMF.LXMessage.SENT, LXMF.LXMessage.PROPAGATED, True),
(LXMF.LXMessage.DELIVERED, LXMF.LXMessage.DIRECT, True),
(LXMF.LXMessage.FAILED, LXMF.LXMessage.DIRECT, True),
(LXMF.LXMessage.CANCELLED, LXMF.LXMessage.DIRECT, True),
):
assert (
is_lxmf_outbound_progress_terminal(_msg(state=state, method=method))
is expected
)
pairs += 1
assert pairs == 6
print("LXMF_SEND_DELIVERY_ORACLE_PROVED")

View file

@ -1114,6 +1114,44 @@ describe("ConversationViewer.vue", () => {
expect(addedItem.is_outbound).toBe(true);
});
it("absorb response then created event do not double-insert the same outbound hash", async () => {
const wrapper = mountConversationViewer();
await flushPromises();
const lxmfMessage = {
hash: "same-outbound-hash",
source_hash: "my-hash",
destination_hash: "test-hash",
content: "once",
state: "sent",
method: "opportunistic",
timestamp: 1700000000,
fields: {},
};
wrapper.vm._absorbOutboundSendResponse({ pendingHash: "pending-x", destinationHash: "test-hash" }, lxmfMessage);
wrapper.vm.onLxmfMessageCreated(lxmfMessage);
const matches = wrapper.vm.chatItems.filter((i) => i.lxmf_message?.hash === "same-outbound-hash");
expect(matches).toHaveLength(1);
});
it("created event then absorb response still keep a single outbound bubble", async () => {
const wrapper = mountConversationViewer();
await flushPromises();
const lxmfMessage = {
hash: "same-outbound-hash-2",
source_hash: "my-hash",
destination_hash: "test-hash",
content: "once",
state: "sent",
method: "opportunistic",
timestamp: 1700000000,
fields: {},
};
wrapper.vm.onLxmfMessageCreated(lxmfMessage);
wrapper.vm._absorbOutboundSendResponse({ pendingHash: "pending-y", destinationHash: "test-hash" }, lxmfMessage);
const matches = wrapper.vm.chatItems.filter((i) => i.lxmf_message?.hash === "same-outbound-hash-2");
expect(matches).toHaveLength(1);
});
it("preserves unknown state for incoming messages", async () => {
const wrapper = mountConversationViewer();

View file

@ -64,9 +64,8 @@ describe("follow-up oracles from exploratory hunt", () => {
});
it("mapViewStateKey scopes TileCache map state by identity", async () => {
const { mapViewStateKey, LEGACY_MAP_STATE_KEY } = await import(
"../../meshchatx/src/frontend/js/mapStateKeys.js"
);
const { mapViewStateKey, LEGACY_MAP_STATE_KEY } =
await import("../../meshchatx/src/frontend/js/mapStateKeys.js");
const a = "aa".repeat(16);
const b = "bb".repeat(16);
expect(mapViewStateKey(a)).not.toBe(mapViewStateKey(b));

View file

@ -124,6 +124,22 @@ describe("outboundMessageStatus LXMF oracle", () => {
expect(isOpportunisticDeferredDelivery(null)).toBe(false);
});
it("sent without delivered is not a hard-fail badge for direct or opportunistic", () => {
for (const method of ["direct", "opportunistic"]) {
const sent = { state: "sent", method };
expect(isOpportunisticDeferredDelivery(sent)).toBe(false);
expect(outboundBubbleStatusIconName(sent)).toBe("check");
expect(outboundBubbleStatusIconName(sent)).not.toBe("check-all");
}
});
it("no-receipt stuck at sent still shows single-check network-sent, not delivered", () => {
const stuck = { state: "sent", method: "opportunistic" };
expect(outboundBubbleStatusIconName(stuck)).toBe("check");
expect(outboundBubbleStatusTitleKey(stuck)).toBe("messages.outbound_sent_network");
expect(isOpportunisticDeferredDelivery(stuck)).toBe(false);
});
it("LXMF_STATUS_UI_ORACLE_PROVED", () => {
let pairs = 0;
for (const state of API_STATES) {

View file

@ -120,10 +120,7 @@ describe("settings/about/interfaces/identities exploratory oracles", () => {
});
it("AboutPage listens for identity-switched and refreshes backups/snapshots", () => {
const src = readFileSync(
join(process.cwd(), "meshchatx/src/frontend/components/about/AboutPage.vue"),
"utf8"
);
const src = readFileSync(join(process.cwd(), "meshchatx/src/frontend/components/about/AboutPage.vue"), "utf8");
expect(src).toContain('GlobalEmitter.on("identity-switched"');
expect(src).toMatch(/onIdentitySwitched\(\)\s*\{[\s\S]*listSnapshots\(\)/);
expect(src).toMatch(/onIdentitySwitched\(\)\s*\{[\s\S]*listAutoBackups\(\)/);