From 78a086e1ea1e965be5bbbe6b61ed316d060f0b37 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 14:51:47 +0000 Subject: [PATCH] Updates --- .../20241206_0003_005_unique_event_hash.py | 63 +++++++++++++++++++ .../collector/handlers/advertisement.py | 23 +++++++ .../collector/handlers/message.py | 22 +++++++ .../collector/handlers/telemetry.py | 23 +++++++ src/meshcore_hub/collector/handlers/trace.py | 22 +++++++ .../common/models/advertisement.py | 6 +- src/meshcore_hub/common/models/message.py | 2 +- src/meshcore_hub/common/models/telemetry.py | 6 +- src/meshcore_hub/common/models/trace_path.py | 2 +- 9 files changed, 159 insertions(+), 10 deletions(-) create mode 100644 alembic/versions/20241206_0003_005_unique_event_hash.py diff --git a/alembic/versions/20241206_0003_005_unique_event_hash.py b/alembic/versions/20241206_0003_005_unique_event_hash.py new file mode 100644 index 0000000..fbcc873 --- /dev/null +++ b/alembic/versions/20241206_0003_005_unique_event_hash.py @@ -0,0 +1,63 @@ +"""Make event_hash columns unique for race condition prevention + +Revision ID: 005 +Revises: 004 +Create Date: 2024-12-06 + +""" + +from typing import Sequence, Union + +from alembic import op + + +# revision identifiers, used by Alembic. +revision: str = "005" +down_revision: Union[str, None] = "004" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + # Drop existing non-unique indexes and create unique constraints + # Note: SQLite handles NULL values as unique (each NULL is distinct) + + # Messages: drop index, create unique constraint + op.drop_index("ix_messages_event_hash", table_name="messages") + op.create_unique_constraint("uq_messages_event_hash", "messages", ["event_hash"]) + + # Advertisements: drop index, create unique constraint + op.drop_index("ix_advertisements_event_hash", table_name="advertisements") + op.create_unique_constraint( + "uq_advertisements_event_hash", "advertisements", ["event_hash"] + ) + + # Trace paths: drop index, create unique constraint + op.drop_index("ix_trace_paths_event_hash", table_name="trace_paths") + op.create_unique_constraint( + "uq_trace_paths_event_hash", "trace_paths", ["event_hash"] + ) + + # Telemetry: drop index, create unique constraint + op.drop_index("ix_telemetry_event_hash", table_name="telemetry") + op.create_unique_constraint("uq_telemetry_event_hash", "telemetry", ["event_hash"]) + + +def downgrade() -> None: + # Restore non-unique indexes + + # Telemetry + op.drop_constraint("uq_telemetry_event_hash", "telemetry", type_="unique") + op.create_index("ix_telemetry_event_hash", "telemetry", ["event_hash"]) + + # Trace paths + op.drop_constraint("uq_trace_paths_event_hash", "trace_paths", type_="unique") + op.create_index("ix_trace_paths_event_hash", "trace_paths", ["event_hash"]) + + # Advertisements + op.drop_constraint("uq_advertisements_event_hash", "advertisements", type_="unique") + op.create_index("ix_advertisements_event_hash", "advertisements", ["event_hash"]) + + # Messages + op.drop_constraint("uq_messages_event_hash", "messages", type_="unique") + op.create_index("ix_messages_event_hash", "messages", ["event_hash"]) diff --git a/src/meshcore_hub/collector/handlers/advertisement.py b/src/meshcore_hub/collector/handlers/advertisement.py index ccab99e..bdf8adb 100644 --- a/src/meshcore_hub/collector/handlers/advertisement.py +++ b/src/meshcore_hub/collector/handlers/advertisement.py @@ -5,6 +5,7 @@ from datetime import datetime, timezone from typing import Any from sqlalchemy import select +from sqlalchemy.exc import IntegrityError from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_advertisement_hash @@ -147,6 +148,28 @@ def handle_advertisement( received_at=now, ) + # Flush to check for duplicate constraint violation (race condition) + try: + session.flush() + except IntegrityError: + # Race condition: another request inserted the same event_hash + session.rollback() + logger.debug( + f"Duplicate advertisement skipped (race condition, " + f"hash={event_hash[:8]}...)" + ) + # Re-add receiver to existing event in a new transaction + if receiver_node: + add_event_receiver( + session=session, + event_type="advertisement", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=None, + received_at=now, + ) + return + logger.info( f"Stored advertisement from {name or adv_public_key[:12]!r} " f"(type={adv_type})" diff --git a/src/meshcore_hub/collector/handlers/message.py b/src/meshcore_hub/collector/handlers/message.py index 1e14343..95105df 100644 --- a/src/meshcore_hub/collector/handlers/message.py +++ b/src/meshcore_hub/collector/handlers/message.py @@ -5,6 +5,7 @@ from datetime import datetime, timezone from typing import Any from sqlalchemy import select +from sqlalchemy.exc import IntegrityError from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_message_hash @@ -163,6 +164,27 @@ def _handle_message( received_at=now, ) + # Flush to check for duplicate constraint violation (race condition) + try: + session.flush() + except IntegrityError: + # Race condition: another request inserted the same event_hash + session.rollback() + logger.debug( + f"Duplicate message skipped (race condition, hash={event_hash[:8]}...)" + ) + # Re-add receiver to existing event in a new transaction + if receiver_node: + add_event_receiver( + session=session, + event_type="message", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=snr, + received_at=now, + ) + return + if message_type == "contact": logger.info( f"Stored contact message from {pubkey_prefix!r}: " diff --git a/src/meshcore_hub/collector/handlers/telemetry.py b/src/meshcore_hub/collector/handlers/telemetry.py index d4e16ca..086135a 100644 --- a/src/meshcore_hub/collector/handlers/telemetry.py +++ b/src/meshcore_hub/collector/handlers/telemetry.py @@ -5,6 +5,7 @@ from datetime import datetime, timezone from typing import Any from sqlalchemy import select +from sqlalchemy.exc import IntegrityError from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_telemetry_hash @@ -138,6 +139,28 @@ def handle_telemetry( received_at=now, ) + # Flush to check for duplicate constraint violation (race condition) + try: + session.flush() + except IntegrityError: + # Race condition: another request inserted the same event_hash + session.rollback() + logger.debug( + f"Duplicate telemetry skipped (race condition, " + f"node={node_public_key[:12]}...)" + ) + # Re-add receiver to existing event in a new transaction + if receiver_node: + add_event_receiver( + session=session, + event_type="telemetry", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=None, + received_at=now, + ) + return + # Log telemetry values if parsed_data: values = ", ".join(f"{k}={v}" for k, v in parsed_data.items()) diff --git a/src/meshcore_hub/collector/handlers/trace.py b/src/meshcore_hub/collector/handlers/trace.py index 3d53d5b..e283c6a 100644 --- a/src/meshcore_hub/collector/handlers/trace.py +++ b/src/meshcore_hub/collector/handlers/trace.py @@ -5,6 +5,7 @@ from datetime import datetime, timezone from typing import Any from sqlalchemy import select +from sqlalchemy.exc import IntegrityError from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_trace_hash @@ -111,4 +112,25 @@ def handle_trace_data( received_at=now, ) + # Flush to check for duplicate constraint violation (race condition) + try: + session.flush() + except IntegrityError: + # Race condition: another request inserted the same event_hash + session.rollback() + logger.debug( + f"Duplicate trace skipped (race condition, tag={initiator_tag})" + ) + # Re-add receiver to existing event in a new transaction + if receiver_node: + add_event_receiver( + session=session, + event_type="trace", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=None, + received_at=now, + ) + return + logger.info(f"Stored trace data: tag={initiator_tag}, hops={hop_count}") diff --git a/src/meshcore_hub/common/models/advertisement.py b/src/meshcore_hub/common/models/advertisement.py index 9fee8c3..0ca55eb 100644 --- a/src/meshcore_hub/common/models/advertisement.py +++ b/src/meshcore_hub/common/models/advertisement.py @@ -61,12 +61,10 @@ class Advertisement(Base, UUIDMixin, TimestampMixin): event_hash: Mapped[Optional[str]] = mapped_column( String(32), nullable=True, + unique=True, ) - __table_args__ = ( - Index("ix_advertisements_received_at", "received_at"), - Index("ix_advertisements_event_hash", "event_hash"), - ) + __table_args__ = (Index("ix_advertisements_received_at", "received_at"),) def __repr__(self) -> str: return f"" diff --git a/src/meshcore_hub/common/models/message.py b/src/meshcore_hub/common/models/message.py index b666ccc..6d0ae1f 100644 --- a/src/meshcore_hub/common/models/message.py +++ b/src/meshcore_hub/common/models/message.py @@ -79,6 +79,7 @@ class Message(Base, UUIDMixin, TimestampMixin): event_hash: Mapped[Optional[str]] = mapped_column( String(32), nullable=True, + unique=True, ) __table_args__ = ( @@ -86,7 +87,6 @@ class Message(Base, UUIDMixin, TimestampMixin): Index("ix_messages_pubkey_prefix", "pubkey_prefix"), Index("ix_messages_channel_idx", "channel_idx"), Index("ix_messages_received_at", "received_at"), - Index("ix_messages_event_hash", "event_hash"), ) def __repr__(self) -> str: diff --git a/src/meshcore_hub/common/models/telemetry.py b/src/meshcore_hub/common/models/telemetry.py index 0d30824..7205e0d 100644 --- a/src/meshcore_hub/common/models/telemetry.py +++ b/src/meshcore_hub/common/models/telemetry.py @@ -57,12 +57,10 @@ class Telemetry(Base, UUIDMixin, TimestampMixin): event_hash: Mapped[Optional[str]] = mapped_column( String(32), nullable=True, + unique=True, ) - __table_args__ = ( - Index("ix_telemetry_received_at", "received_at"), - Index("ix_telemetry_event_hash", "event_hash"), - ) + __table_args__ = (Index("ix_telemetry_received_at", "received_at"),) def __repr__(self) -> str: return ( diff --git a/src/meshcore_hub/common/models/trace_path.py b/src/meshcore_hub/common/models/trace_path.py index 51f9582..caf1132 100644 --- a/src/meshcore_hub/common/models/trace_path.py +++ b/src/meshcore_hub/common/models/trace_path.py @@ -70,12 +70,12 @@ class TracePath(Base, UUIDMixin, TimestampMixin): event_hash: Mapped[Optional[str]] = mapped_column( String(32), nullable=True, + unique=True, ) __table_args__ = ( Index("ix_trace_paths_initiator_tag", "initiator_tag"), Index("ix_trace_paths_received_at", "received_at"), - Index("ix_trace_paths_event_hash", "event_hash"), ) def __repr__(self) -> str: