From c80986fe67d0b531070d690201429398d6f7d1e5 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 6 Dec 2025 12:08:07 +0000 Subject: [PATCH 01/10] Add event deduplication at presentation layer When multiple receiver nodes are running, the same mesh events (messages, advertisements) are reported multiple times. This causes duplicate entries in the Web UI. Changes: - Add hash_utils.py with deterministic hash functions for each event type - Add `dedupe` parameter to messages and advertisements API endpoints (default: True) - Update dashboard stats to use distinct counts for messages/advertisements - Deduplicate recent advertisements and channel messages in dashboard - Add comprehensive tests for hash utilities Hash strategy: - Messages: hash of text + pubkey_prefix + channel_idx + sender_timestamp + txt_type - Advertisements: hash of public_key + name + adv_type + flags + 5-minute time bucket --- src/meshcore_hub/api/routes/advertisements.py | 64 ++++- src/meshcore_hub/api/routes/dashboard.py | 141 +++++++-- src/meshcore_hub/api/routes/messages.py | 59 +++- src/meshcore_hub/common/hash_utils.py | 144 ++++++++++ tests/test_common/test_hash_utils.py | 268 ++++++++++++++++++ 5 files changed, 644 insertions(+), 32 deletions(-) create mode 100644 src/meshcore_hub/common/hash_utils.py create mode 100644 tests/test_common/test_hash_utils.py diff --git a/src/meshcore_hub/api/routes/advertisements.py b/src/meshcore_hub/api/routes/advertisements.py index 55a57ea..94e74ff 100644 --- a/src/meshcore_hub/api/routes/advertisements.py +++ b/src/meshcore_hub/api/routes/advertisements.py @@ -9,6 +9,7 @@ from sqlalchemy.orm import aliased, selectinload from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession +from meshcore_hub.common.hash_utils import compute_advertisement_hash from meshcore_hub.common.models import Advertisement, Node from meshcore_hub.common.schemas.messages import AdvertisementList, AdvertisementRead @@ -35,6 +36,9 @@ async def list_advertisements( ), since: Optional[datetime] = Query(None, description="Start timestamp"), until: Optional[datetime] = Query(None, description="End timestamp"), + dedupe: bool = Query( + True, description="Deduplicate advertisements from multiple receivers" + ), limit: int = Query(50, ge=1, le=100, description="Page size"), offset: int = Query(0, ge=0, description="Page offset"), ) -> AdvertisementList: @@ -70,12 +74,37 @@ async def list_advertisements( if until: query = query.where(Advertisement.received_at <= until) - # Get total count - count_query = select(func.count()).select_from(query.subquery()) - total = session.execute(count_query).scalar() or 0 + # When deduplicating, we need to fetch more results and compute distinct count + if dedupe: + # For deduplicated count, count distinct by public_key within time buckets + # We use a 5-minute time bucket for advertisements + distinct_subquery = ( + select( + Advertisement.public_key, + Advertisement.name, + Advertisement.adv_type, + Advertisement.flags, + # Use date truncation for time bucketing (5 min = 300 seconds) + (func.strftime("%s", Advertisement.received_at) / 300).label( + "time_bucket" + ), + ) + .distinct() + .select_from(query.subquery()) + ) + count_query = select(func.count()).select_from(distinct_subquery.subquery()) + total = session.execute(count_query).scalar() or 0 - # Apply pagination - query = query.order_by(Advertisement.received_at.desc()).offset(offset).limit(limit) + # Fetch extra results to account for duplicates + fetch_limit = (limit + offset) * 3 + query = query.order_by(Advertisement.received_at.desc()).limit(fetch_limit) + else: + # Standard count and pagination + count_query = select(func.count()).select_from(query.subquery()) + total = session.execute(count_query).scalar() or 0 + query = ( + query.order_by(Advertisement.received_at.desc()).offset(offset).limit(limit) + ) # Execute results = session.execute(query).all() @@ -99,8 +128,25 @@ async def list_advertisements( # Build response with node details items = [] + seen_hashes: set[str] = set() + for row in results: adv = row[0] + + # Compute hash for deduplication + if dedupe: + adv_hash = compute_advertisement_hash( + public_key=adv.public_key, + name=adv.name, + adv_type=adv.adv_type, + flags=adv.flags, + received_at=adv.received_at, + bucket_minutes=5, + ) + if adv_hash in seen_hashes: + continue + seen_hashes.add(adv_hash) + receiver_node = nodes_by_id.get(row.receiver_id) if row.receiver_id else None source_node = nodes_by_id.get(row.source_id) if row.source_id else None @@ -119,6 +165,14 @@ async def list_advertisements( } items.append(AdvertisementRead(**data)) + # Stop once we have enough items (for dedupe mode with pagination) + if dedupe and len(items) >= offset + limit: + break + + # Apply offset for dedupe mode (we fetched from beginning) + if dedupe: + items = items[offset : offset + limit] + return AdvertisementList( items=items, total=total, diff --git a/src/meshcore_hub/api/routes/dashboard.py b/src/meshcore_hub/api/routes/dashboard.py index 9afecdb..d2bebf0 100644 --- a/src/meshcore_hub/api/routes/dashboard.py +++ b/src/meshcore_hub/api/routes/dashboard.py @@ -8,6 +8,10 @@ from sqlalchemy import func, select from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession +from meshcore_hub.common.hash_utils import ( + compute_advertisement_hash, + compute_message_hash, +) from meshcore_hub.common.models import Advertisement, Message, Node, NodeTag from meshcore_hub.common.schemas.messages import ( ChannelMessage, @@ -43,45 +47,110 @@ async def get_stats( or 0 ) - # Total messages + # Total messages (deduplicated by content hash) + distinct_messages = ( + select( + Message.text, + Message.pubkey_prefix, + Message.channel_idx, + Message.sender_timestamp, + Message.txt_type, + ) + .distinct() + .subquery() + ) total_messages = ( - session.execute(select(func.count()).select_from(Message)).scalar() or 0 + session.execute(select(func.count()).select_from(distinct_messages)).scalar() + or 0 ) - # Messages today + # Messages today (deduplicated) + distinct_messages_today = ( + select( + Message.text, + Message.pubkey_prefix, + Message.channel_idx, + Message.sender_timestamp, + Message.txt_type, + ) + .where(Message.received_at >= today_start) + .distinct() + .subquery() + ) messages_today = ( session.execute( - select(func.count()) - .select_from(Message) - .where(Message.received_at >= today_start) + select(func.count()).select_from(distinct_messages_today) ).scalar() or 0 ) - # Total advertisements + # Total advertisements (deduplicated by public_key + 5min time bucket) + distinct_advertisements = ( + select( + Advertisement.public_key, + Advertisement.name, + Advertisement.adv_type, + Advertisement.flags, + (func.strftime("%s", Advertisement.received_at) / 300).label("time_bucket"), + ) + .distinct() + .subquery() + ) total_advertisements = ( - session.execute(select(func.count()).select_from(Advertisement)).scalar() or 0 + session.execute( + select(func.count()).select_from(distinct_advertisements) + ).scalar() + or 0 ) - # Advertisements in last 24h + # Advertisements in last 24h (deduplicated) + distinct_advertisements_24h = ( + select( + Advertisement.public_key, + Advertisement.name, + Advertisement.adv_type, + Advertisement.flags, + (func.strftime("%s", Advertisement.received_at) / 300).label("time_bucket"), + ) + .where(Advertisement.received_at >= yesterday) + .distinct() + .subquery() + ) advertisements_24h = ( session.execute( - select(func.count()) - .select_from(Advertisement) - .where(Advertisement.received_at >= yesterday) + select(func.count()).select_from(distinct_advertisements_24h) ).scalar() or 0 ) - # Recent advertisements (last 10) - recent_ads = ( + # Recent advertisements (last 10, deduplicated) + # Fetch more to ensure we have 10 unique after deduplication + recent_ads_raw = ( session.execute( - select(Advertisement).order_by(Advertisement.received_at.desc()).limit(10) + select(Advertisement).order_by(Advertisement.received_at.desc()).limit(30) ) .scalars() .all() ) + # Deduplicate by hash + seen_ad_hashes: set[str] = set() + recent_ads = [] + for ad in recent_ads_raw: + ad_hash = compute_advertisement_hash( + public_key=ad.public_key, + name=ad.name, + adv_type=ad.adv_type, + flags=ad.flags, + received_at=ad.received_at, + bucket_minutes=5, + ) + if ad_hash not in seen_ad_hashes: + seen_ad_hashes.add(ad_hash) + recent_ads.append(ad) + if len(recent_ads) >= 10: + break + # Get node names, adv_types, and friendly_name tags for the advertised nodes ad_public_keys = [ad.public_key for ad in recent_ads] node_names: dict[str, str] = {} @@ -119,29 +188,57 @@ async def get_stats( for ad in recent_ads ] - # Channel message counts - channel_counts_query = ( - select(Message.channel_idx, func.count()) + # Channel message counts (deduplicated) + distinct_channel_messages = ( + select( + Message.channel_idx, + Message.text, + Message.pubkey_prefix, + Message.sender_timestamp, + Message.txt_type, + ) .where(Message.message_type == "channel") .where(Message.channel_idx.isnot(None)) - .group_by(Message.channel_idx) + .distinct() + .subquery() ) + channel_counts_query = select( + distinct_channel_messages.c.channel_idx, func.count() + ).group_by(distinct_channel_messages.c.channel_idx) channel_results = session.execute(channel_counts_query).all() channel_message_counts = { int(channel): int(count) for channel, count in channel_results } - # Get latest 5 messages for each channel that has messages + # Get latest 5 messages for each channel that has messages (deduplicated) channel_messages: dict[int, list[ChannelMessage]] = {} for channel_idx, _ in channel_results: + # Fetch more messages to deduplicate messages_query = ( select(Message) .where(Message.message_type == "channel") .where(Message.channel_idx == channel_idx) .order_by(Message.received_at.desc()) - .limit(5) + .limit(15) ) - channel_msgs = session.execute(messages_query).scalars().all() + channel_msgs_raw = session.execute(messages_query).scalars().all() + + # Deduplicate + seen_msg_hashes: set[str] = set() + channel_msgs = [] + for m in channel_msgs_raw: + msg_hash = compute_message_hash( + text=m.text, + pubkey_prefix=m.pubkey_prefix, + channel_idx=m.channel_idx, + sender_timestamp=m.sender_timestamp, + txt_type=m.txt_type, + ) + if msg_hash not in seen_msg_hashes: + seen_msg_hashes.add(msg_hash) + channel_msgs.append(m) + if len(channel_msgs) >= 5: + break # Look up sender names for these messages msg_prefixes = [m.pubkey_prefix for m in channel_msgs if m.pubkey_prefix] diff --git a/src/meshcore_hub/api/routes/messages.py b/src/meshcore_hub/api/routes/messages.py index 79a461a..f9ee382 100644 --- a/src/meshcore_hub/api/routes/messages.py +++ b/src/meshcore_hub/api/routes/messages.py @@ -9,6 +9,7 @@ from sqlalchemy.orm import aliased, selectinload from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession +from meshcore_hub.common.hash_utils import compute_message_hash from meshcore_hub.common.models import Message, Node, NodeTag from meshcore_hub.common.schemas.messages import MessageList, MessageRead @@ -38,6 +39,9 @@ async def list_messages( since: Optional[datetime] = Query(None, description="Start timestamp"), until: Optional[datetime] = Query(None, description="End timestamp"), search: Optional[str] = Query(None, description="Search in message text"), + dedupe: bool = Query( + True, description="Deduplicate messages from multiple receivers" + ), limit: int = Query(50, ge=1, le=100, description="Page size"), offset: int = Query(0, ge=0, description="Page offset"), ) -> MessageList: @@ -74,12 +78,33 @@ async def list_messages( if search: query = query.where(Message.text.ilike(f"%{search}%")) - # Get total count - count_query = select(func.count()).select_from(query.subquery()) - total = session.execute(count_query).scalar() or 0 + # When deduplicating, we need to fetch more results to ensure we have enough + # after removing duplicates, and compute distinct count differently + if dedupe: + # For deduplicated count, we need to count distinct content + # Use a subquery that groups by content-identifying fields + distinct_subquery = ( + select( + Message.text, + Message.pubkey_prefix, + Message.channel_idx, + Message.sender_timestamp, + Message.txt_type, + ) + .distinct() + .select_from(query.subquery()) + ) + count_query = select(func.count()).select_from(distinct_subquery.subquery()) + total = session.execute(count_query).scalar() or 0 - # Apply pagination - query = query.order_by(Message.received_at.desc()).offset(offset).limit(limit) + # Fetch extra results to account for duplicates (3x limit + offset) + fetch_limit = (limit + offset) * 3 + query = query.order_by(Message.received_at.desc()).limit(fetch_limit) + else: + # Standard count and pagination + count_query = select(func.count()).select_from(query.subquery()) + total = session.execute(count_query).scalar() or 0 + query = query.order_by(Message.received_at.desc()).offset(offset).limit(limit) # Execute results = session.execute(query).all() @@ -128,8 +153,24 @@ async def list_messages( # Build response with sender info and received_by items = [] + seen_hashes: set[str] = set() + for row in results: m = row[0] + + # Compute hash for deduplication + if dedupe: + msg_hash = compute_message_hash( + text=m.text, + pubkey_prefix=m.pubkey_prefix, + channel_idx=m.channel_idx, + sender_timestamp=m.sender_timestamp, + txt_type=m.txt_type, + ) + if msg_hash in seen_hashes: + continue + seen_hashes.add(msg_hash) + receiver_pk = row.receiver_pk receiver_name = row.receiver_name receiver_node = ( @@ -162,6 +203,14 @@ async def list_messages( } items.append(MessageRead(**msg_dict)) + # Stop once we have enough items (for dedupe mode with pagination) + if dedupe and len(items) >= offset + limit: + break + + # Apply offset for dedupe mode (we fetched from beginning) + if dedupe: + items = items[offset : offset + limit] + return MessageList( items=items, total=total, diff --git a/src/meshcore_hub/common/hash_utils.py b/src/meshcore_hub/common/hash_utils.py new file mode 100644 index 0000000..5be1825 --- /dev/null +++ b/src/meshcore_hub/common/hash_utils.py @@ -0,0 +1,144 @@ +"""Event hash utilities for deduplication. + +This module provides functions to compute deterministic hashes for events, +allowing deduplication when multiple receiver nodes report the same event. +""" + +import hashlib +from datetime import datetime +from typing import Optional + + +def compute_message_hash( + text: str, + pubkey_prefix: Optional[str] = None, + channel_idx: Optional[int] = None, + sender_timestamp: Optional[datetime] = None, + txt_type: Optional[int] = None, +) -> str: + """Compute a deterministic hash for a message. + + The hash is computed from fields that uniquely identify a message's content + and sender, excluding receiver-specific data. + + Args: + text: Message content + pubkey_prefix: Sender's public key prefix (12 chars) + channel_idx: Channel index for channel messages + sender_timestamp: Sender's timestamp + txt_type: Message type indicator + + Returns: + 32-character hex hash string + """ + # Build a canonical string from the relevant fields + parts = [ + text or "", + pubkey_prefix or "", + str(channel_idx) if channel_idx is not None else "", + sender_timestamp.isoformat() if sender_timestamp else "", + str(txt_type) if txt_type is not None else "", + ] + canonical = "|".join(parts) + return hashlib.md5(canonical.encode("utf-8")).hexdigest() + + +def compute_advertisement_hash( + public_key: str, + name: Optional[str] = None, + adv_type: Optional[str] = None, + flags: Optional[int] = None, + received_at: Optional[datetime] = None, + bucket_minutes: int = 5, +) -> str: + """Compute a deterministic hash for an advertisement. + + Advertisements are bucketed by time since the same node may advertise + periodically and we want to deduplicate within a time window. + + Args: + public_key: Advertised node's public key + name: Advertised name + adv_type: Node type + flags: Capability flags + received_at: When received (used for time bucketing) + bucket_minutes: Time bucket size in minutes (default 5) + + Returns: + 32-character hex hash string + """ + # Bucket the time to allow deduplication within a window + time_bucket = "" + if received_at: + # Round down to nearest bucket + bucket_seconds = bucket_minutes * 60 + epoch = int(received_at.timestamp()) + bucket_epoch = (epoch // bucket_seconds) * bucket_seconds + time_bucket = str(bucket_epoch) + + parts = [ + public_key, + name or "", + adv_type or "", + str(flags) if flags is not None else "", + time_bucket, + ] + canonical = "|".join(parts) + return hashlib.md5(canonical.encode("utf-8")).hexdigest() + + +def compute_trace_hash(initiator_tag: int) -> str: + """Compute a deterministic hash for a trace path. + + Trace paths have a unique initiator_tag that serves as the identifier. + + Args: + initiator_tag: Unique trace identifier + + Returns: + 32-character hex hash string + """ + return hashlib.md5(str(initiator_tag).encode("utf-8")).hexdigest() + + +def compute_telemetry_hash( + node_public_key: str, + parsed_data: Optional[dict] = None, + received_at: Optional[datetime] = None, + bucket_minutes: int = 5, +) -> str: + """Compute a deterministic hash for a telemetry record. + + Telemetry is bucketed by time since nodes report periodically. + + Args: + node_public_key: Reporting node's public key + parsed_data: Decoded sensor readings + received_at: When received (used for time bucketing) + bucket_minutes: Time bucket size in minutes (default 5) + + Returns: + 32-character hex hash string + """ + # Bucket the time + time_bucket = "" + if received_at: + bucket_seconds = bucket_minutes * 60 + epoch = int(received_at.timestamp()) + bucket_epoch = (epoch // bucket_seconds) * bucket_seconds + time_bucket = str(bucket_epoch) + + # Serialize parsed_data deterministically + data_str = "" + if parsed_data: + # Sort keys for deterministic serialization + sorted_items = sorted(parsed_data.items()) + data_str = str(sorted_items) + + parts = [ + node_public_key, + data_str, + time_bucket, + ] + canonical = "|".join(parts) + return hashlib.md5(canonical.encode("utf-8")).hexdigest() diff --git a/tests/test_common/test_hash_utils.py b/tests/test_common/test_hash_utils.py new file mode 100644 index 0000000..6ba9e48 --- /dev/null +++ b/tests/test_common/test_hash_utils.py @@ -0,0 +1,268 @@ +"""Tests for hash utilities for event deduplication.""" + +from datetime import datetime, timezone + +import pytest + +from meshcore_hub.common.hash_utils import ( + compute_advertisement_hash, + compute_message_hash, + compute_telemetry_hash, + compute_trace_hash, +) + + +class TestComputeMessageHash: + """Tests for compute_message_hash function.""" + + def test_same_content_produces_same_hash(self) -> None: + """Identical messages should produce the same hash.""" + timestamp = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + + hash1 = compute_message_hash( + text="Hello World", + pubkey_prefix="01ab2186c4d5", + channel_idx=4, + sender_timestamp=timestamp, + txt_type=1, + ) + hash2 = compute_message_hash( + text="Hello World", + pubkey_prefix="01ab2186c4d5", + channel_idx=4, + sender_timestamp=timestamp, + txt_type=1, + ) + + assert hash1 == hash2 + + def test_different_text_produces_different_hash(self) -> None: + """Messages with different text should have different hashes.""" + timestamp = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + + hash1 = compute_message_hash( + text="Hello World", + pubkey_prefix="01ab2186c4d5", + sender_timestamp=timestamp, + ) + hash2 = compute_message_hash( + text="Goodbye World", + pubkey_prefix="01ab2186c4d5", + sender_timestamp=timestamp, + ) + + assert hash1 != hash2 + + def test_different_sender_produces_different_hash(self) -> None: + """Messages from different senders should have different hashes.""" + timestamp = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + + hash1 = compute_message_hash( + text="Hello", + pubkey_prefix="01ab2186c4d5", + sender_timestamp=timestamp, + ) + hash2 = compute_message_hash( + text="Hello", + pubkey_prefix="99ff8877aabb", + sender_timestamp=timestamp, + ) + + assert hash1 != hash2 + + def test_different_channel_produces_different_hash(self) -> None: + """Messages on different channels should have different hashes.""" + hash1 = compute_message_hash(text="Hello", channel_idx=1) + hash2 = compute_message_hash(text="Hello", channel_idx=2) + + assert hash1 != hash2 + + def test_handles_none_values(self) -> None: + """Hash function should handle None values gracefully.""" + hash1 = compute_message_hash( + text="Test", + pubkey_prefix=None, + channel_idx=None, + sender_timestamp=None, + txt_type=None, + ) + + assert hash1 is not None + assert len(hash1) == 32 # MD5 hex digest length + + +class TestComputeAdvertisementHash: + """Tests for compute_advertisement_hash function.""" + + def test_same_content_same_bucket_produces_same_hash(self) -> None: + """Advertisements within the same time bucket should match.""" + # Two times within the same 5-minute bucket + time1 = datetime(2024, 1, 15, 10, 31, 0, tzinfo=timezone.utc) + time2 = datetime(2024, 1, 15, 10, 33, 0, tzinfo=timezone.utc) + + hash1 = compute_advertisement_hash( + public_key="a" * 64, + name="Node1", + adv_type="chat", + flags=128, + received_at=time1, + bucket_minutes=5, + ) + hash2 = compute_advertisement_hash( + public_key="a" * 64, + name="Node1", + adv_type="chat", + flags=128, + received_at=time2, + bucket_minutes=5, + ) + + assert hash1 == hash2 + + def test_different_bucket_produces_different_hash(self) -> None: + """Advertisements in different time buckets should not match.""" + # Two times in different 5-minute buckets + time1 = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + time2 = datetime(2024, 1, 15, 10, 36, 0, tzinfo=timezone.utc) + + hash1 = compute_advertisement_hash( + public_key="a" * 64, + name="Node1", + received_at=time1, + bucket_minutes=5, + ) + hash2 = compute_advertisement_hash( + public_key="a" * 64, + name="Node1", + received_at=time2, + bucket_minutes=5, + ) + + assert hash1 != hash2 + + def test_different_public_key_produces_different_hash(self) -> None: + """Advertisements from different nodes should have different hashes.""" + time = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + + hash1 = compute_advertisement_hash( + public_key="a" * 64, + received_at=time, + ) + hash2 = compute_advertisement_hash( + public_key="b" * 64, + received_at=time, + ) + + assert hash1 != hash2 + + def test_configurable_bucket_size(self) -> None: + """Bucket size should be configurable.""" + time1 = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + time2 = datetime(2024, 1, 15, 10, 35, 0, tzinfo=timezone.utc) + + # With 5-minute bucket, these should be in different buckets + hash1_5min = compute_advertisement_hash( + public_key="a" * 64, + received_at=time1, + bucket_minutes=5, + ) + hash2_5min = compute_advertisement_hash( + public_key="a" * 64, + received_at=time2, + bucket_minutes=5, + ) + assert hash1_5min != hash2_5min + + # With 10-minute bucket, these should be in the same bucket + hash1_10min = compute_advertisement_hash( + public_key="a" * 64, + received_at=time1, + bucket_minutes=10, + ) + hash2_10min = compute_advertisement_hash( + public_key="a" * 64, + received_at=time2, + bucket_minutes=10, + ) + assert hash1_10min == hash2_10min + + +class TestComputeTraceHash: + """Tests for compute_trace_hash function.""" + + def test_same_tag_produces_same_hash(self) -> None: + """Same initiator_tag should produce same hash.""" + hash1 = compute_trace_hash(initiator_tag=123456789) + hash2 = compute_trace_hash(initiator_tag=123456789) + + assert hash1 == hash2 + + def test_different_tag_produces_different_hash(self) -> None: + """Different initiator_tag should produce different hash.""" + hash1 = compute_trace_hash(initiator_tag=123456789) + hash2 = compute_trace_hash(initiator_tag=987654321) + + assert hash1 != hash2 + + +class TestComputeTelemetryHash: + """Tests for compute_telemetry_hash function.""" + + def test_same_content_same_bucket_produces_same_hash(self) -> None: + """Telemetry within the same time bucket should match.""" + time1 = datetime(2024, 1, 15, 10, 31, 0, tzinfo=timezone.utc) + time2 = datetime(2024, 1, 15, 10, 33, 0, tzinfo=timezone.utc) + data = {"temperature": 22.5, "humidity": 65} + + hash1 = compute_telemetry_hash( + node_public_key="a" * 64, + parsed_data=data, + received_at=time1, + bucket_minutes=5, + ) + hash2 = compute_telemetry_hash( + node_public_key="a" * 64, + parsed_data=data, + received_at=time2, + bucket_minutes=5, + ) + + assert hash1 == hash2 + + def test_different_data_produces_different_hash(self) -> None: + """Different sensor readings should produce different hashes.""" + time = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + + hash1 = compute_telemetry_hash( + node_public_key="a" * 64, + parsed_data={"temperature": 22.5}, + received_at=time, + ) + hash2 = compute_telemetry_hash( + node_public_key="a" * 64, + parsed_data={"temperature": 25.0}, + received_at=time, + ) + + assert hash1 != hash2 + + def test_deterministic_dict_serialization(self) -> None: + """Dict serialization should be deterministic regardless of key order.""" + time = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) + + # Same data, different key order in source dicts + data1 = {"a": 1, "b": 2, "c": 3} + data2 = {"c": 3, "a": 1, "b": 2} + + hash1 = compute_telemetry_hash( + node_public_key="a" * 64, + parsed_data=data1, + received_at=time, + ) + hash2 = compute_telemetry_hash( + node_public_key="a" * 64, + parsed_data=data2, + received_at=time, + ) + + assert hash1 == hash2 From 6e3b86a1ad823fe22fc7985d2e650933e1ba1034 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 6 Dec 2025 12:23:14 +0000 Subject: [PATCH 02/10] Add collector-level event deduplication using content hashes Replace presentation-layer deduplication with collector-level approach: - Add event_hash column to messages, advertisements, trace_paths, telemetry tables - Handlers compute content hashes and skip duplicate events at insertion time - Use 5-minute time buckets for advertisements and telemetry - Include Alembic migration for schema changes --- .../20241206_0001_003_add_event_hash.py | 67 +++++++++ src/meshcore_hub/api/routes/advertisements.py | 64 +------- src/meshcore_hub/api/routes/dashboard.py | 141 +++--------------- src/meshcore_hub/api/routes/messages.py | 59 +------- .../collector/handlers/advertisement.py | 26 ++++ .../collector/handlers/message.py | 20 +++ .../collector/handlers/telemetry.py | 21 +++ src/meshcore_hub/collector/handlers/trace.py | 14 ++ .../common/models/advertisement.py | 9 +- src/meshcore_hub/common/models/message.py | 5 + src/meshcore_hub/common/models/telemetry.py | 9 +- src/meshcore_hub/common/models/trace_path.py | 7 +- tests/test_common/test_hash_utils.py | 2 - 13 files changed, 207 insertions(+), 237 deletions(-) create mode 100644 alembic/versions/20241206_0001_003_add_event_hash.py diff --git a/alembic/versions/20241206_0001_003_add_event_hash.py b/alembic/versions/20241206_0001_003_add_event_hash.py new file mode 100644 index 0000000..085f036 --- /dev/null +++ b/alembic/versions/20241206_0001_003_add_event_hash.py @@ -0,0 +1,67 @@ +"""Add event_hash column to event tables for deduplication + +Revision ID: 003 +Revises: 002 +Create Date: 2024-12-06 + +""" + +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = "003" +down_revision: Union[str, None] = "002" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + # Add event_hash column to messages table + op.add_column( + "messages", + sa.Column("event_hash", sa.String(32), nullable=True), + ) + op.create_index("ix_messages_event_hash", "messages", ["event_hash"]) + + # Add event_hash column to advertisements table + op.add_column( + "advertisements", + sa.Column("event_hash", sa.String(32), nullable=True), + ) + op.create_index("ix_advertisements_event_hash", "advertisements", ["event_hash"]) + + # Add event_hash column to trace_paths table + op.add_column( + "trace_paths", + sa.Column("event_hash", sa.String(32), nullable=True), + ) + op.create_index("ix_trace_paths_event_hash", "trace_paths", ["event_hash"]) + + # Add event_hash column to telemetry table + op.add_column( + "telemetry", + sa.Column("event_hash", sa.String(32), nullable=True), + ) + op.create_index("ix_telemetry_event_hash", "telemetry", ["event_hash"]) + + +def downgrade() -> None: + # Remove event_hash from telemetry + op.drop_index("ix_telemetry_event_hash", table_name="telemetry") + op.drop_column("telemetry", "event_hash") + + # Remove event_hash from trace_paths + op.drop_index("ix_trace_paths_event_hash", table_name="trace_paths") + op.drop_column("trace_paths", "event_hash") + + # Remove event_hash from advertisements + op.drop_index("ix_advertisements_event_hash", table_name="advertisements") + op.drop_column("advertisements", "event_hash") + + # Remove event_hash from messages + op.drop_index("ix_messages_event_hash", table_name="messages") + op.drop_column("messages", "event_hash") diff --git a/src/meshcore_hub/api/routes/advertisements.py b/src/meshcore_hub/api/routes/advertisements.py index 94e74ff..55a57ea 100644 --- a/src/meshcore_hub/api/routes/advertisements.py +++ b/src/meshcore_hub/api/routes/advertisements.py @@ -9,7 +9,6 @@ from sqlalchemy.orm import aliased, selectinload from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession -from meshcore_hub.common.hash_utils import compute_advertisement_hash from meshcore_hub.common.models import Advertisement, Node from meshcore_hub.common.schemas.messages import AdvertisementList, AdvertisementRead @@ -36,9 +35,6 @@ async def list_advertisements( ), since: Optional[datetime] = Query(None, description="Start timestamp"), until: Optional[datetime] = Query(None, description="End timestamp"), - dedupe: bool = Query( - True, description="Deduplicate advertisements from multiple receivers" - ), limit: int = Query(50, ge=1, le=100, description="Page size"), offset: int = Query(0, ge=0, description="Page offset"), ) -> AdvertisementList: @@ -74,37 +70,12 @@ async def list_advertisements( if until: query = query.where(Advertisement.received_at <= until) - # When deduplicating, we need to fetch more results and compute distinct count - if dedupe: - # For deduplicated count, count distinct by public_key within time buckets - # We use a 5-minute time bucket for advertisements - distinct_subquery = ( - select( - Advertisement.public_key, - Advertisement.name, - Advertisement.adv_type, - Advertisement.flags, - # Use date truncation for time bucketing (5 min = 300 seconds) - (func.strftime("%s", Advertisement.received_at) / 300).label( - "time_bucket" - ), - ) - .distinct() - .select_from(query.subquery()) - ) - count_query = select(func.count()).select_from(distinct_subquery.subquery()) - total = session.execute(count_query).scalar() or 0 + # Get total count + count_query = select(func.count()).select_from(query.subquery()) + total = session.execute(count_query).scalar() or 0 - # Fetch extra results to account for duplicates - fetch_limit = (limit + offset) * 3 - query = query.order_by(Advertisement.received_at.desc()).limit(fetch_limit) - else: - # Standard count and pagination - count_query = select(func.count()).select_from(query.subquery()) - total = session.execute(count_query).scalar() or 0 - query = ( - query.order_by(Advertisement.received_at.desc()).offset(offset).limit(limit) - ) + # Apply pagination + query = query.order_by(Advertisement.received_at.desc()).offset(offset).limit(limit) # Execute results = session.execute(query).all() @@ -128,25 +99,8 @@ async def list_advertisements( # Build response with node details items = [] - seen_hashes: set[str] = set() - for row in results: adv = row[0] - - # Compute hash for deduplication - if dedupe: - adv_hash = compute_advertisement_hash( - public_key=adv.public_key, - name=adv.name, - adv_type=adv.adv_type, - flags=adv.flags, - received_at=adv.received_at, - bucket_minutes=5, - ) - if adv_hash in seen_hashes: - continue - seen_hashes.add(adv_hash) - receiver_node = nodes_by_id.get(row.receiver_id) if row.receiver_id else None source_node = nodes_by_id.get(row.source_id) if row.source_id else None @@ -165,14 +119,6 @@ async def list_advertisements( } items.append(AdvertisementRead(**data)) - # Stop once we have enough items (for dedupe mode with pagination) - if dedupe and len(items) >= offset + limit: - break - - # Apply offset for dedupe mode (we fetched from beginning) - if dedupe: - items = items[offset : offset + limit] - return AdvertisementList( items=items, total=total, diff --git a/src/meshcore_hub/api/routes/dashboard.py b/src/meshcore_hub/api/routes/dashboard.py index d2bebf0..9afecdb 100644 --- a/src/meshcore_hub/api/routes/dashboard.py +++ b/src/meshcore_hub/api/routes/dashboard.py @@ -8,10 +8,6 @@ from sqlalchemy import func, select from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession -from meshcore_hub.common.hash_utils import ( - compute_advertisement_hash, - compute_message_hash, -) from meshcore_hub.common.models import Advertisement, Message, Node, NodeTag from meshcore_hub.common.schemas.messages import ( ChannelMessage, @@ -47,110 +43,45 @@ async def get_stats( or 0 ) - # Total messages (deduplicated by content hash) - distinct_messages = ( - select( - Message.text, - Message.pubkey_prefix, - Message.channel_idx, - Message.sender_timestamp, - Message.txt_type, - ) - .distinct() - .subquery() - ) + # Total messages total_messages = ( - session.execute(select(func.count()).select_from(distinct_messages)).scalar() - or 0 + session.execute(select(func.count()).select_from(Message)).scalar() or 0 ) - # Messages today (deduplicated) - distinct_messages_today = ( - select( - Message.text, - Message.pubkey_prefix, - Message.channel_idx, - Message.sender_timestamp, - Message.txt_type, - ) - .where(Message.received_at >= today_start) - .distinct() - .subquery() - ) + # Messages today messages_today = ( session.execute( - select(func.count()).select_from(distinct_messages_today) + select(func.count()) + .select_from(Message) + .where(Message.received_at >= today_start) ).scalar() or 0 ) - # Total advertisements (deduplicated by public_key + 5min time bucket) - distinct_advertisements = ( - select( - Advertisement.public_key, - Advertisement.name, - Advertisement.adv_type, - Advertisement.flags, - (func.strftime("%s", Advertisement.received_at) / 300).label("time_bucket"), - ) - .distinct() - .subquery() - ) + # Total advertisements total_advertisements = ( - session.execute( - select(func.count()).select_from(distinct_advertisements) - ).scalar() - or 0 + session.execute(select(func.count()).select_from(Advertisement)).scalar() or 0 ) - # Advertisements in last 24h (deduplicated) - distinct_advertisements_24h = ( - select( - Advertisement.public_key, - Advertisement.name, - Advertisement.adv_type, - Advertisement.flags, - (func.strftime("%s", Advertisement.received_at) / 300).label("time_bucket"), - ) - .where(Advertisement.received_at >= yesterday) - .distinct() - .subquery() - ) + # Advertisements in last 24h advertisements_24h = ( session.execute( - select(func.count()).select_from(distinct_advertisements_24h) + select(func.count()) + .select_from(Advertisement) + .where(Advertisement.received_at >= yesterday) ).scalar() or 0 ) - # Recent advertisements (last 10, deduplicated) - # Fetch more to ensure we have 10 unique after deduplication - recent_ads_raw = ( + # Recent advertisements (last 10) + recent_ads = ( session.execute( - select(Advertisement).order_by(Advertisement.received_at.desc()).limit(30) + select(Advertisement).order_by(Advertisement.received_at.desc()).limit(10) ) .scalars() .all() ) - # Deduplicate by hash - seen_ad_hashes: set[str] = set() - recent_ads = [] - for ad in recent_ads_raw: - ad_hash = compute_advertisement_hash( - public_key=ad.public_key, - name=ad.name, - adv_type=ad.adv_type, - flags=ad.flags, - received_at=ad.received_at, - bucket_minutes=5, - ) - if ad_hash not in seen_ad_hashes: - seen_ad_hashes.add(ad_hash) - recent_ads.append(ad) - if len(recent_ads) >= 10: - break - # Get node names, adv_types, and friendly_name tags for the advertised nodes ad_public_keys = [ad.public_key for ad in recent_ads] node_names: dict[str, str] = {} @@ -188,57 +119,29 @@ async def get_stats( for ad in recent_ads ] - # Channel message counts (deduplicated) - distinct_channel_messages = ( - select( - Message.channel_idx, - Message.text, - Message.pubkey_prefix, - Message.sender_timestamp, - Message.txt_type, - ) + # Channel message counts + channel_counts_query = ( + select(Message.channel_idx, func.count()) .where(Message.message_type == "channel") .where(Message.channel_idx.isnot(None)) - .distinct() - .subquery() + .group_by(Message.channel_idx) ) - channel_counts_query = select( - distinct_channel_messages.c.channel_idx, func.count() - ).group_by(distinct_channel_messages.c.channel_idx) channel_results = session.execute(channel_counts_query).all() channel_message_counts = { int(channel): int(count) for channel, count in channel_results } - # Get latest 5 messages for each channel that has messages (deduplicated) + # Get latest 5 messages for each channel that has messages channel_messages: dict[int, list[ChannelMessage]] = {} for channel_idx, _ in channel_results: - # Fetch more messages to deduplicate messages_query = ( select(Message) .where(Message.message_type == "channel") .where(Message.channel_idx == channel_idx) .order_by(Message.received_at.desc()) - .limit(15) + .limit(5) ) - channel_msgs_raw = session.execute(messages_query).scalars().all() - - # Deduplicate - seen_msg_hashes: set[str] = set() - channel_msgs = [] - for m in channel_msgs_raw: - msg_hash = compute_message_hash( - text=m.text, - pubkey_prefix=m.pubkey_prefix, - channel_idx=m.channel_idx, - sender_timestamp=m.sender_timestamp, - txt_type=m.txt_type, - ) - if msg_hash not in seen_msg_hashes: - seen_msg_hashes.add(msg_hash) - channel_msgs.append(m) - if len(channel_msgs) >= 5: - break + channel_msgs = session.execute(messages_query).scalars().all() # Look up sender names for these messages msg_prefixes = [m.pubkey_prefix for m in channel_msgs if m.pubkey_prefix] diff --git a/src/meshcore_hub/api/routes/messages.py b/src/meshcore_hub/api/routes/messages.py index f9ee382..79a461a 100644 --- a/src/meshcore_hub/api/routes/messages.py +++ b/src/meshcore_hub/api/routes/messages.py @@ -9,7 +9,6 @@ from sqlalchemy.orm import aliased, selectinload from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession -from meshcore_hub.common.hash_utils import compute_message_hash from meshcore_hub.common.models import Message, Node, NodeTag from meshcore_hub.common.schemas.messages import MessageList, MessageRead @@ -39,9 +38,6 @@ async def list_messages( since: Optional[datetime] = Query(None, description="Start timestamp"), until: Optional[datetime] = Query(None, description="End timestamp"), search: Optional[str] = Query(None, description="Search in message text"), - dedupe: bool = Query( - True, description="Deduplicate messages from multiple receivers" - ), limit: int = Query(50, ge=1, le=100, description="Page size"), offset: int = Query(0, ge=0, description="Page offset"), ) -> MessageList: @@ -78,33 +74,12 @@ async def list_messages( if search: query = query.where(Message.text.ilike(f"%{search}%")) - # When deduplicating, we need to fetch more results to ensure we have enough - # after removing duplicates, and compute distinct count differently - if dedupe: - # For deduplicated count, we need to count distinct content - # Use a subquery that groups by content-identifying fields - distinct_subquery = ( - select( - Message.text, - Message.pubkey_prefix, - Message.channel_idx, - Message.sender_timestamp, - Message.txt_type, - ) - .distinct() - .select_from(query.subquery()) - ) - count_query = select(func.count()).select_from(distinct_subquery.subquery()) - total = session.execute(count_query).scalar() or 0 + # Get total count + count_query = select(func.count()).select_from(query.subquery()) + total = session.execute(count_query).scalar() or 0 - # Fetch extra results to account for duplicates (3x limit + offset) - fetch_limit = (limit + offset) * 3 - query = query.order_by(Message.received_at.desc()).limit(fetch_limit) - else: - # Standard count and pagination - count_query = select(func.count()).select_from(query.subquery()) - total = session.execute(count_query).scalar() or 0 - query = query.order_by(Message.received_at.desc()).offset(offset).limit(limit) + # Apply pagination + query = query.order_by(Message.received_at.desc()).offset(offset).limit(limit) # Execute results = session.execute(query).all() @@ -153,24 +128,8 @@ async def list_messages( # Build response with sender info and received_by items = [] - seen_hashes: set[str] = set() - for row in results: m = row[0] - - # Compute hash for deduplication - if dedupe: - msg_hash = compute_message_hash( - text=m.text, - pubkey_prefix=m.pubkey_prefix, - channel_idx=m.channel_idx, - sender_timestamp=m.sender_timestamp, - txt_type=m.txt_type, - ) - if msg_hash in seen_hashes: - continue - seen_hashes.add(msg_hash) - receiver_pk = row.receiver_pk receiver_name = row.receiver_name receiver_node = ( @@ -203,14 +162,6 @@ async def list_messages( } items.append(MessageRead(**msg_dict)) - # Stop once we have enough items (for dedupe mode with pagination) - if dedupe and len(items) >= offset + limit: - break - - # Apply offset for dedupe mode (we fetched from beginning) - if dedupe: - items = items[offset : offset + limit] - return MessageList( items=items, total=total, diff --git a/src/meshcore_hub/collector/handlers/advertisement.py b/src/meshcore_hub/collector/handlers/advertisement.py index 1588a61..b9d6fbf 100644 --- a/src/meshcore_hub/collector/handlers/advertisement.py +++ b/src/meshcore_hub/collector/handlers/advertisement.py @@ -7,6 +7,7 @@ from typing import Any from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager +from meshcore_hub.common.hash_utils import compute_advertisement_hash from meshcore_hub.common.models import Advertisement, Node logger = logging.getLogger(__name__) @@ -40,7 +41,31 @@ def handle_advertisement( flags = payload.get("flags") now = datetime.now(timezone.utc) + # Compute event hash for deduplication (5-minute time bucket) + event_hash = compute_advertisement_hash( + public_key=adv_public_key, + name=name, + adv_type=adv_type, + flags=flags, + received_at=now, + bucket_minutes=5, + ) + with db.session_scope() as session: + # Check if advertisement with same hash already exists + existing = session.execute( + select(Advertisement.id).where(Advertisement.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + logger.debug(f"Duplicate advertisement skipped (hash={event_hash[:8]}...)") + # Still update node last_seen even for duplicate advertisements + node_query = select(Node).where(Node.public_key == adv_public_key) + node = session.execute(node_query).scalar_one_or_none() + if node: + node.last_seen = now + return + # Find or create receiver node receiver_node = None if public_key: @@ -91,6 +116,7 @@ def handle_advertisement( adv_type=adv_type, flags=flags, received_at=now, + event_hash=event_hash, ) session.add(advertisement) diff --git a/src/meshcore_hub/collector/handlers/message.py b/src/meshcore_hub/collector/handlers/message.py index 2494cac..47a54f6 100644 --- a/src/meshcore_hub/collector/handlers/message.py +++ b/src/meshcore_hub/collector/handlers/message.py @@ -7,6 +7,7 @@ from typing import Any from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager +from meshcore_hub.common.hash_utils import compute_message_hash from meshcore_hub.common.models import Message, Node logger = logging.getLogger(__name__) @@ -84,7 +85,25 @@ def _handle_message( except (ValueError, OSError): pass + # Compute event hash for deduplication + event_hash = compute_message_hash( + text=text, + pubkey_prefix=pubkey_prefix, + channel_idx=channel_idx, + sender_timestamp=sender_timestamp, + txt_type=txt_type, + ) + with db.session_scope() as session: + # Check if message with same hash already exists + existing = session.execute( + select(Message.id).where(Message.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + logger.debug(f"Duplicate message skipped (hash={event_hash[:8]}...)") + return + # Find receiver node receiver_node = None if public_key: @@ -115,6 +134,7 @@ def _handle_message( snr=snr, sender_timestamp=sender_timestamp, received_at=now, + event_hash=event_hash, ) session.add(message) diff --git a/src/meshcore_hub/collector/handlers/telemetry.py b/src/meshcore_hub/collector/handlers/telemetry.py index 40c14e1..6e0bfe9 100644 --- a/src/meshcore_hub/collector/handlers/telemetry.py +++ b/src/meshcore_hub/collector/handlers/telemetry.py @@ -7,6 +7,7 @@ from typing import Any from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager +from meshcore_hub.common.hash_utils import compute_telemetry_hash from meshcore_hub.common.models import Node, Telemetry logger = logging.getLogger(__name__) @@ -49,7 +50,26 @@ def handle_telemetry( except ValueError: lpp_bytes = lpp_data.encode() + # Compute event hash for deduplication (5-minute time bucket) + event_hash = compute_telemetry_hash( + node_public_key=node_public_key, + parsed_data=parsed_data, + received_at=now, + bucket_minutes=5, + ) + with db.session_scope() as session: + # Check if telemetry with same hash already exists + existing = session.execute( + select(Telemetry.id).where(Telemetry.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + logger.debug( + f"Duplicate telemetry skipped (node={node_public_key[:12]}...)" + ) + return + # Find receiver node receiver_node = None if public_key: @@ -92,6 +112,7 @@ def handle_telemetry( lpp_data=lpp_bytes, parsed_data=parsed_data, received_at=now, + event_hash=event_hash, ) session.add(telemetry) diff --git a/src/meshcore_hub/collector/handlers/trace.py b/src/meshcore_hub/collector/handlers/trace.py index 646860d..280f764 100644 --- a/src/meshcore_hub/collector/handlers/trace.py +++ b/src/meshcore_hub/collector/handlers/trace.py @@ -7,6 +7,7 @@ from typing import Any from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager +from meshcore_hub.common.hash_utils import compute_trace_hash from meshcore_hub.common.models import Node, TracePath logger = logging.getLogger(__name__) @@ -40,7 +41,19 @@ def handle_trace_data( snr_values = payload.get("snr_values") hop_count = payload.get("hop_count") + # Compute event hash for deduplication (initiator_tag is unique per trace) + event_hash = compute_trace_hash(initiator_tag=initiator_tag) + with db.session_scope() as session: + # Check if trace with same hash already exists + existing = session.execute( + select(TracePath.id).where(TracePath.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + logger.debug(f"Duplicate trace skipped (tag={initiator_tag})") + return + # Find receiver node receiver_node = None if public_key: @@ -69,6 +82,7 @@ def handle_trace_data( snr_values=snr_values, hop_count=hop_count, received_at=now, + event_hash=event_hash, ) session.add(trace_path) diff --git a/src/meshcore_hub/common/models/advertisement.py b/src/meshcore_hub/common/models/advertisement.py index 192cd94..9fee8c3 100644 --- a/src/meshcore_hub/common/models/advertisement.py +++ b/src/meshcore_hub/common/models/advertisement.py @@ -58,8 +58,15 @@ class Advertisement(Base, UUIDMixin, TimestampMixin): default=utc_now, nullable=False, ) + event_hash: Mapped[Optional[str]] = mapped_column( + String(32), + nullable=True, + ) - __table_args__ = (Index("ix_advertisements_received_at", "received_at"),) + __table_args__ = ( + Index("ix_advertisements_received_at", "received_at"), + Index("ix_advertisements_event_hash", "event_hash"), + ) 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 ad50630..b666ccc 100644 --- a/src/meshcore_hub/common/models/message.py +++ b/src/meshcore_hub/common/models/message.py @@ -76,12 +76,17 @@ class Message(Base, UUIDMixin, TimestampMixin): default=utc_now, nullable=False, ) + event_hash: Mapped[Optional[str]] = mapped_column( + String(32), + nullable=True, + ) __table_args__ = ( Index("ix_messages_message_type", "message_type"), 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 77084ac..0d30824 100644 --- a/src/meshcore_hub/common/models/telemetry.py +++ b/src/meshcore_hub/common/models/telemetry.py @@ -54,8 +54,15 @@ class Telemetry(Base, UUIDMixin, TimestampMixin): default=utc_now, nullable=False, ) + event_hash: Mapped[Optional[str]] = mapped_column( + String(32), + nullable=True, + ) - __table_args__ = (Index("ix_telemetry_received_at", "received_at"),) + __table_args__ = ( + Index("ix_telemetry_received_at", "received_at"), + Index("ix_telemetry_event_hash", "event_hash"), + ) 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 62c45b3..51f9582 100644 --- a/src/meshcore_hub/common/models/trace_path.py +++ b/src/meshcore_hub/common/models/trace_path.py @@ -3,7 +3,7 @@ from datetime import datetime from typing import Optional -from sqlalchemy import BigInteger, DateTime, ForeignKey, Index, Integer +from sqlalchemy import BigInteger, DateTime, ForeignKey, Index, Integer, String from sqlalchemy.dialects.sqlite import JSON from sqlalchemy.orm import Mapped, mapped_column @@ -67,10 +67,15 @@ class TracePath(Base, UUIDMixin, TimestampMixin): default=utc_now, nullable=False, ) + event_hash: Mapped[Optional[str]] = mapped_column( + String(32), + nullable=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: diff --git a/tests/test_common/test_hash_utils.py b/tests/test_common/test_hash_utils.py index 6ba9e48..1ca94c4 100644 --- a/tests/test_common/test_hash_utils.py +++ b/tests/test_common/test_hash_utils.py @@ -2,8 +2,6 @@ from datetime import datetime, timezone -import pytest - from meshcore_hub.common.hash_utils import ( compute_advertisement_hash, compute_message_hash, From 2ea04deb7e63042c3bddd769cfab7f4260bfcf4d Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 12:53:29 +0000 Subject: [PATCH 03/10] Updates --- README.md | 19 +++++------ src/meshcore_hub/__main__.py | 34 +++++++++++++++++++ src/meshcore_hub/api/app.py | 3 +- src/meshcore_hub/collector/cli.py | 12 +++---- .../collector/handlers/advertisement.py | 3 +- .../collector/handlers/telemetry.py | 3 +- src/meshcore_hub/collector/subscriber.py | 10 +++--- src/meshcore_hub/common/hash_utils.py | 10 +++--- 8 files changed, 59 insertions(+), 35 deletions(-) diff --git a/README.md b/README.md index 60cc1bb..f8e4d29 100644 --- a/README.md +++ b/README.md @@ -180,24 +180,21 @@ cd meshcore-hub cp .env.example .env # Edit .env with your settings (API keys, serial port, network info) -# Option 1: Start core services (mqtt, collector, api, web) +# Create database schema +docker compose --profile migrate run --rm db-migrate + +# Seed the database +docker compose --profile seed run --rm seed + +# Start core services (mqtt, collector, api, web) docker compose up -d -# Option 2: Start with mock device for testing -docker compose --profile mock up -d - -# Option 3: Start with real MeshCore device +# Start sender/receiver interface docker compose --profile interface-receiver up -d # View logs docker compose logs -f -# Run database migrations (one-time) -docker compose --profile migrate up - -# Import seed data manually (also runs on collector startup) -docker compose --profile seed up - # Stop services docker compose down ``` diff --git a/src/meshcore_hub/__main__.py b/src/meshcore_hub/__main__.py index 96d8fc9..871615b 100644 --- a/src/meshcore_hub/__main__.py +++ b/src/meshcore_hub/__main__.py @@ -174,6 +174,40 @@ def db_history() -> None: command.history(alembic_cfg) +@db.command("stamp") +@click.option( + "--revision", + type=str, + default="head", + help="Target revision to stamp (default: head)", +) +@click.option( + "--database-url", + type=str, + default=None, + envvar="DATABASE_URL", + help="Database connection URL", +) +def db_stamp(revision: str, database_url: str | None) -> None: + """Stamp database with revision without running migrations. + + Use this to mark an existing database as up-to-date when the schema + was created before Alembic migrations were introduced. + """ + import os + from alembic import command + from alembic.config import Config + + click.echo(f"Stamping database with revision: {revision}") + + alembic_cfg = Config("alembic.ini") + if database_url: + os.environ["DATABASE_URL"] = database_url + + command.stamp(alembic_cfg, revision) + click.echo("Database stamped successfully.") + + # Health check commands for Docker HEALTHCHECK @cli.group() def health() -> None: diff --git a/src/meshcore_hub/api/app.py b/src/meshcore_hub/api/app.py index 4a1efd1..fa9191d 100644 --- a/src/meshcore_hub/api/app.py +++ b/src/meshcore_hub/api/app.py @@ -32,10 +32,9 @@ async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]: # Get database URL from app state database_url = getattr(app.state, "database_url", "sqlite:///./meshcore.db") - # Initialize database + # Initialize database (schema managed by Alembic migrations) logger.info(f"Initializing database: {database_url}") _db_manager = DatabaseManager(database_url) - _db_manager.create_tables() yield diff --git a/src/meshcore_hub/collector/cli.py b/src/meshcore_hub/collector/cli.py index c501de9..655d56f 100644 --- a/src/meshcore_hub/collector/cli.py +++ b/src/meshcore_hub/collector/cli.py @@ -182,11 +182,10 @@ def _run_collector_service( click.echo(f"MQTT: {mqtt_host}:{mqtt_port} (prefix: {prefix})") click.echo(f"Database: {database_url}") - # Initialize database and run seed import on startup + # Initialize database (schema managed by Alembic migrations) from meshcore_hub.common.database import DatabaseManager db = DatabaseManager(database_url) - db.create_tables() # Auto-seed from seed files on startup click.echo("") @@ -294,9 +293,8 @@ def seed_cmd( from meshcore_hub.common.database import DatabaseManager - # Initialize database + # Initialize database (schema managed by Alembic migrations) db = DatabaseManager(ctx.obj["database_url"]) - db.create_tables() # Run seed import imported_any = _run_seed_import( @@ -448,9 +446,8 @@ def import_tags_cmd( from meshcore_hub.common.database import DatabaseManager from meshcore_hub.collector.tag_import import import_tags - # Initialize database + # Initialize database (schema managed by Alembic migrations) db = DatabaseManager(ctx.obj["database_url"]) - db.create_tables() # Import tags stats = import_tags( @@ -529,9 +526,8 @@ def import_members_cmd( from meshcore_hub.common.database import DatabaseManager from meshcore_hub.collector.member_import import import_members - # Initialize database + # Initialize database (schema managed by Alembic migrations) db = DatabaseManager(ctx.obj["database_url"]) - db.create_tables() # Import members stats = import_members( diff --git a/src/meshcore_hub/collector/handlers/advertisement.py b/src/meshcore_hub/collector/handlers/advertisement.py index b9d6fbf..8bb04f9 100644 --- a/src/meshcore_hub/collector/handlers/advertisement.py +++ b/src/meshcore_hub/collector/handlers/advertisement.py @@ -41,14 +41,13 @@ def handle_advertisement( flags = payload.get("flags") now = datetime.now(timezone.utc) - # Compute event hash for deduplication (5-minute time bucket) + # Compute event hash for deduplication (30-second time bucket) event_hash = compute_advertisement_hash( public_key=adv_public_key, name=name, adv_type=adv_type, flags=flags, received_at=now, - bucket_minutes=5, ) with db.session_scope() as session: diff --git a/src/meshcore_hub/collector/handlers/telemetry.py b/src/meshcore_hub/collector/handlers/telemetry.py index 6e0bfe9..61ebbf4 100644 --- a/src/meshcore_hub/collector/handlers/telemetry.py +++ b/src/meshcore_hub/collector/handlers/telemetry.py @@ -50,12 +50,11 @@ def handle_telemetry( except ValueError: lpp_bytes = lpp_data.encode() - # Compute event hash for deduplication (5-minute time bucket) + # Compute event hash for deduplication (30-second time bucket) event_hash = compute_telemetry_hash( node_public_key=node_public_key, parsed_data=parsed_data, received_at=now, - bucket_minutes=5, ) with db.session_scope() as session: diff --git a/src/meshcore_hub/collector/subscriber.py b/src/meshcore_hub/collector/subscriber.py index 57d7ccf..0b4e4e2 100644 --- a/src/meshcore_hub/collector/subscriber.py +++ b/src/meshcore_hub/collector/subscriber.py @@ -206,14 +206,16 @@ class Subscriber: """Start the subscriber.""" logger.info("Starting collector subscriber") - # Create database tables if needed + # Verify database connection (schema managed by Alembic migrations) try: - self.db.create_tables() + # Test connection by getting a session + session = self.db.get_session() + session.close() self._db_connected = True - logger.info("Database initialized") + logger.info("Database connection verified") except Exception as e: self._db_connected = False - logger.error(f"Failed to initialize database: {e}") + logger.error(f"Failed to connect to database: {e}") raise # Connect to MQTT broker diff --git a/src/meshcore_hub/common/hash_utils.py b/src/meshcore_hub/common/hash_utils.py index 5be1825..5e75a9a 100644 --- a/src/meshcore_hub/common/hash_utils.py +++ b/src/meshcore_hub/common/hash_utils.py @@ -49,7 +49,7 @@ def compute_advertisement_hash( adv_type: Optional[str] = None, flags: Optional[int] = None, received_at: Optional[datetime] = None, - bucket_minutes: int = 5, + bucket_seconds: int = 30, ) -> str: """Compute a deterministic hash for an advertisement. @@ -62,7 +62,7 @@ def compute_advertisement_hash( adv_type: Node type flags: Capability flags received_at: When received (used for time bucketing) - bucket_minutes: Time bucket size in minutes (default 5) + bucket_seconds: Time bucket size in seconds (default 30) Returns: 32-character hex hash string @@ -71,7 +71,6 @@ def compute_advertisement_hash( time_bucket = "" if received_at: # Round down to nearest bucket - bucket_seconds = bucket_minutes * 60 epoch = int(received_at.timestamp()) bucket_epoch = (epoch // bucket_seconds) * bucket_seconds time_bucket = str(bucket_epoch) @@ -105,7 +104,7 @@ def compute_telemetry_hash( node_public_key: str, parsed_data: Optional[dict] = None, received_at: Optional[datetime] = None, - bucket_minutes: int = 5, + bucket_seconds: int = 30, ) -> str: """Compute a deterministic hash for a telemetry record. @@ -115,7 +114,7 @@ def compute_telemetry_hash( node_public_key: Reporting node's public key parsed_data: Decoded sensor readings received_at: When received (used for time bucketing) - bucket_minutes: Time bucket size in minutes (default 5) + bucket_seconds: Time bucket size in seconds (default 30) Returns: 32-character hex hash string @@ -123,7 +122,6 @@ def compute_telemetry_hash( # Bucket the time time_bucket = "" if received_at: - bucket_seconds = bucket_minutes * 60 epoch = int(received_at.timestamp()) bucket_epoch = (epoch // bucket_seconds) * bucket_seconds time_bucket = str(bucket_epoch) From d715e4e4f0ffff2d2885b68ed163cc7fde637095 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 13:33:02 +0000 Subject: [PATCH 04/10] Updates --- .../20241206_0002_004_event_receivers.py | 63 +++++++++ src/meshcore_hub/api/routes/advertisements.py | 82 ++++++++++- src/meshcore_hub/api/routes/messages.py | 90 ++++++++++++- .../collector/handlers/advertisement.py | 60 ++++++--- .../collector/handlers/message.py | 47 +++++-- .../collector/handlers/telemetry.py | 49 +++++-- src/meshcore_hub/collector/handlers/trace.py | 47 +++++-- src/meshcore_hub/common/models/__init__.py | 3 + .../common/models/event_receiver.py | 127 ++++++++++++++++++ src/meshcore_hub/common/schemas/__init__.py | 4 +- src/meshcore_hub/common/schemas/messages.py | 33 +++++ .../web/templates/advertisements.html | 26 +++- src/meshcore_hub/web/templates/messages.html | 29 +++- tests/test_common/test_hash_utils.py | 28 ++-- 14 files changed, 616 insertions(+), 72 deletions(-) create mode 100644 alembic/versions/20241206_0002_004_event_receivers.py create mode 100644 src/meshcore_hub/common/models/event_receiver.py diff --git a/alembic/versions/20241206_0002_004_event_receivers.py b/alembic/versions/20241206_0002_004_event_receivers.py new file mode 100644 index 0000000..2682803 --- /dev/null +++ b/alembic/versions/20241206_0002_004_event_receivers.py @@ -0,0 +1,63 @@ +"""Add event_receivers junction table for multi-receiver tracking + +Revision ID: 004 +Revises: 003 +Create Date: 2024-12-06 + +""" + +from typing import Sequence, Union + +from alembic import op +import sqlalchemy as sa + + +# revision identifiers, used by Alembic. +revision: str = "004" +down_revision: Union[str, None] = "003" +branch_labels: Union[str, Sequence[str], None] = None +depends_on: Union[str, Sequence[str], None] = None + + +def upgrade() -> None: + op.create_table( + "event_receivers", + sa.Column("id", sa.String(36), primary_key=True), + sa.Column("event_type", sa.String(20), nullable=False), + sa.Column("event_hash", sa.String(32), nullable=False), + sa.Column( + "receiver_node_id", + sa.String(36), + sa.ForeignKey("nodes.id", ondelete="CASCADE"), + nullable=False, + ), + sa.Column("snr", sa.Float, nullable=True), + sa.Column("received_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("updated_at", sa.DateTime(timezone=True), nullable=False), + sa.UniqueConstraint( + "event_hash", "receiver_node_id", name="uq_event_receivers_hash_node" + ), + ) + op.create_index( + "ix_event_receivers_event_hash", + "event_receivers", + ["event_hash"], + ) + op.create_index( + "ix_event_receivers_receiver_node_id", + "event_receivers", + ["receiver_node_id"], + ) + op.create_index( + "ix_event_receivers_type_hash", + "event_receivers", + ["event_type", "event_hash"], + ) + + +def downgrade() -> None: + op.drop_index("ix_event_receivers_type_hash", table_name="event_receivers") + op.drop_index("ix_event_receivers_receiver_node_id", table_name="event_receivers") + op.drop_index("ix_event_receivers_event_hash", table_name="event_receivers") + op.drop_table("event_receivers") diff --git a/src/meshcore_hub/api/routes/advertisements.py b/src/meshcore_hub/api/routes/advertisements.py index 55a57ea..e36974b 100644 --- a/src/meshcore_hub/api/routes/advertisements.py +++ b/src/meshcore_hub/api/routes/advertisements.py @@ -9,8 +9,12 @@ from sqlalchemy.orm import aliased, selectinload from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession -from meshcore_hub.common.models import Advertisement, Node -from meshcore_hub.common.schemas.messages import AdvertisementList, AdvertisementRead +from meshcore_hub.common.models import Advertisement, EventReceiver, Node, NodeTag +from meshcore_hub.common.schemas.messages import ( + AdvertisementList, + AdvertisementRead, + ReceiverInfo, +) router = APIRouter() @@ -25,6 +29,62 @@ def _get_friendly_name(node: Optional[Node]) -> Optional[str]: return None +def _fetch_receivers_for_events( + session: DbSession, + event_type: str, + event_hashes: list[str], +) -> dict[str, list[ReceiverInfo]]: + """Fetch receiver info for a list of events by their hashes.""" + if not event_hashes: + return {} + + query = ( + select( + EventReceiver.event_hash, + EventReceiver.snr, + EventReceiver.received_at, + Node.id.label("node_id"), + Node.public_key, + Node.name, + ) + .join(Node, EventReceiver.receiver_node_id == Node.id) + .where(EventReceiver.event_type == event_type) + .where(EventReceiver.event_hash.in_(event_hashes)) + .order_by(EventReceiver.received_at) + ) + + results = session.execute(query).all() + receivers_by_hash: dict[str, list[ReceiverInfo]] = {} + + node_ids = [r.node_id for r in results] + friendly_names: dict[str, str] = {} + if node_ids: + fn_query = ( + select(NodeTag.node_id, NodeTag.value) + .where(NodeTag.node_id.in_(node_ids)) + .where(NodeTag.key == "friendly_name") + ) + for node_id, value in session.execute(fn_query).all(): + friendly_names[node_id] = value + + for row in results: + if row.event_hash not in receivers_by_hash: + receivers_by_hash[row.event_hash] = [] + + receivers_by_hash[row.event_hash].append( + ReceiverInfo( + node_id=row.node_id, + public_key=row.public_key, + name=row.name, + friendly_name=friendly_names.get(row.node_id), + snr=row.snr, + received_at=row.received_at, + ) + ) + + return receivers_by_hash + + @router.get("", response_model=AdvertisementList) async def list_advertisements( _: RequireRead, @@ -97,6 +157,12 @@ async def list_advertisements( nodes = session.execute(nodes_query).scalars().all() nodes_by_id = {n.id: n for n in nodes} + # Fetch all receivers for these advertisements + event_hashes = [r[0].event_hash for r in results if r[0].event_hash] + receivers_by_hash = _fetch_receivers_for_events( + session, "advertisement", event_hashes + ) + # Build response with node details items = [] for row in results: @@ -116,6 +182,9 @@ async def list_advertisements( "flags": adv.flags, "received_at": adv.received_at, "created_at": adv.created_at, + "receivers": ( + receivers_by_hash.get(adv.event_hash, []) if adv.event_hash else [] + ), } items.append(AdvertisementRead(**data)) @@ -175,6 +244,14 @@ async def get_advertisement( receiver_node = nodes_by_id.get(result.receiver_id) if result.receiver_id else None source_node = nodes_by_id.get(result.source_id) if result.source_id else None + # Fetch receivers for this advertisement + receivers = [] + if adv.event_hash: + receivers_by_hash = _fetch_receivers_for_events( + session, "advertisement", [adv.event_hash] + ) + receivers = receivers_by_hash.get(adv.event_hash, []) + data = { "received_by": result.receiver_pk, "receiver_name": result.receiver_name, @@ -187,5 +264,6 @@ async def get_advertisement( "flags": adv.flags, "received_at": adv.received_at, "created_at": adv.created_at, + "receivers": receivers, } return AdvertisementRead(**data) diff --git a/src/meshcore_hub/api/routes/messages.py b/src/meshcore_hub/api/routes/messages.py index 79a461a..811ad53 100644 --- a/src/meshcore_hub/api/routes/messages.py +++ b/src/meshcore_hub/api/routes/messages.py @@ -9,8 +9,8 @@ from sqlalchemy.orm import aliased, selectinload from meshcore_hub.api.auth import RequireRead from meshcore_hub.api.dependencies import DbSession -from meshcore_hub.common.models import Message, Node, NodeTag -from meshcore_hub.common.schemas.messages import MessageList, MessageRead +from meshcore_hub.common.models import EventReceiver, Message, Node, NodeTag +from meshcore_hub.common.schemas.messages import MessageList, MessageRead, ReceiverInfo router = APIRouter() @@ -25,6 +25,75 @@ def _get_friendly_name(node: Optional[Node]) -> Optional[str]: return None +def _fetch_receivers_for_events( + session: DbSession, + event_type: str, + event_hashes: list[str], +) -> dict[str, list[ReceiverInfo]]: + """Fetch receiver info for a list of events by their hashes. + + Args: + session: Database session + event_type: Type of event ('message', 'advertisement', etc.) + event_hashes: List of event hashes to fetch receivers for + + Returns: + Dict mapping event_hash to list of ReceiverInfo objects + """ + if not event_hashes: + return {} + + # Query event_receivers with receiver node info + query = ( + select( + EventReceiver.event_hash, + EventReceiver.snr, + EventReceiver.received_at, + Node.id.label("node_id"), + Node.public_key, + Node.name, + ) + .join(Node, EventReceiver.receiver_node_id == Node.id) + .where(EventReceiver.event_type == event_type) + .where(EventReceiver.event_hash.in_(event_hashes)) + .order_by(EventReceiver.received_at) + ) + + results = session.execute(query).all() + + # Group by event_hash + receivers_by_hash: dict[str, list[ReceiverInfo]] = {} + + # Get friendly names for receiver nodes + node_ids = [r.node_id for r in results] + friendly_names: dict[str, str] = {} + if node_ids: + fn_query = ( + select(NodeTag.node_id, NodeTag.value) + .where(NodeTag.node_id.in_(node_ids)) + .where(NodeTag.key == "friendly_name") + ) + for node_id, value in session.execute(fn_query).all(): + friendly_names[node_id] = value + + for row in results: + if row.event_hash not in receivers_by_hash: + receivers_by_hash[row.event_hash] = [] + + receivers_by_hash[row.event_hash].append( + ReceiverInfo( + node_id=row.node_id, + public_key=row.public_key, + name=row.name, + friendly_name=friendly_names.get(row.node_id), + snr=row.snr, + received_at=row.received_at, + ) + ) + + return receivers_by_hash + + @router.get("", response_model=MessageList) async def list_messages( _: RequireRead, @@ -126,6 +195,10 @@ async def list_messages( receivers = session.execute(receivers_query).scalars().all() receivers_by_id = {n.id: n for n in receivers} + # Fetch all receivers for these messages + event_hashes = [r[0].event_hash for r in results if r[0].event_hash] + receivers_by_hash = _fetch_receivers_for_events(session, "message", event_hashes) + # Build response with sender info and received_by items = [] for row in results: @@ -159,6 +232,9 @@ async def list_messages( "sender_timestamp": m.sender_timestamp, "received_at": m.received_at, "created_at": m.created_at, + "receivers": ( + receivers_by_hash.get(m.event_hash, []) if m.event_hash else [] + ), } items.append(MessageRead(**msg_dict)) @@ -189,6 +265,15 @@ async def get_message( raise HTTPException(status_code=404, detail="Message not found") message, receiver_pk = result + + # Fetch receivers for this message + receivers = [] + if message.event_hash: + receivers_by_hash = _fetch_receivers_for_events( + session, "message", [message.event_hash] + ) + receivers = receivers_by_hash.get(message.event_hash, []) + data = { "id": message.id, "receiver_node_id": message.receiver_node_id, @@ -204,5 +289,6 @@ async def get_message( "sender_timestamp": message.sender_timestamp, "received_at": message.received_at, "created_at": message.created_at, + "receivers": receivers, } return MessageRead(**data) diff --git a/src/meshcore_hub/collector/handlers/advertisement.py b/src/meshcore_hub/collector/handlers/advertisement.py index 8bb04f9..ccab99e 100644 --- a/src/meshcore_hub/collector/handlers/advertisement.py +++ b/src/meshcore_hub/collector/handlers/advertisement.py @@ -8,7 +8,7 @@ from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_advertisement_hash -from meshcore_hub.common.models import Advertisement, Node +from meshcore_hub.common.models import Advertisement, Node, add_event_receiver logger = logging.getLogger(__name__) @@ -51,21 +51,7 @@ def handle_advertisement( ) with db.session_scope() as session: - # Check if advertisement with same hash already exists - existing = session.execute( - select(Advertisement.id).where(Advertisement.event_hash == event_hash) - ).scalar_one_or_none() - - if existing: - logger.debug(f"Duplicate advertisement skipped (hash={event_hash[:8]}...)") - # Still update node last_seen even for duplicate advertisements - node_query = select(Node).where(Node.public_key == adv_public_key) - node = session.execute(node_query).scalar_one_or_none() - if node: - node.last_seen = now - return - - # Find or create receiver node + # Find or create receiver node first (needed for both new and duplicate events) receiver_node = None if public_key: receiver_query = select(Node).where(Node.public_key == public_key) @@ -79,6 +65,37 @@ def handle_advertisement( ) session.add(receiver_node) session.flush() + else: + receiver_node.last_seen = now + + # Check if advertisement with same hash already exists + existing = session.execute( + select(Advertisement.id).where(Advertisement.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + # Still update advertised node's last_seen even for duplicate advertisements + node_query = select(Node).where(Node.public_key == adv_public_key) + node = session.execute(node_query).scalar_one_or_none() + if node: + node.last_seen = now + + # Add this receiver to the junction table + if receiver_node: + added = add_event_receiver( + session=session, + event_type="advertisement", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=None, # Advertisements don't have SNR + received_at=now, + ) + if added: + logger.debug( + f"Added receiver {public_key[:12]}... to advertisement " + f"(hash={event_hash[:8]}...)" + ) + return # Find or create advertised node node_query = select(Node).where(Node.public_key == adv_public_key) @@ -119,6 +136,17 @@ def handle_advertisement( ) session.add(advertisement) + # Add first receiver to junction table + 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, + ) + 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 47a54f6..1e14343 100644 --- a/src/meshcore_hub/collector/handlers/message.py +++ b/src/meshcore_hub/collector/handlers/message.py @@ -8,7 +8,7 @@ from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_message_hash -from meshcore_hub.common.models import Message, Node +from meshcore_hub.common.models import Message, Node, add_event_receiver logger = logging.getLogger(__name__) @@ -95,16 +95,7 @@ def _handle_message( ) with db.session_scope() as session: - # Check if message with same hash already exists - existing = session.execute( - select(Message.id).where(Message.event_hash == event_hash) - ).scalar_one_or_none() - - if existing: - logger.debug(f"Duplicate message skipped (hash={event_hash[:8]}...)") - return - - # Find receiver node + # Find or create receiver node first (needed for both new and duplicate events) receiver_node = None if public_key: receiver_query = select(Node).where(Node.public_key == public_key) @@ -121,6 +112,29 @@ def _handle_message( else: receiver_node.last_seen = now + # Check if message with same hash already exists + existing = session.execute( + select(Message.id).where(Message.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + # Event already exists - just add this receiver to the junction table + if receiver_node: + added = add_event_receiver( + session=session, + event_type="message", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=snr, + received_at=now, + ) + if added: + logger.debug( + f"Added receiver {public_key[:12]}... to message " + f"(hash={event_hash[:8]}...)" + ) + return + # Create message record message = Message( receiver_node_id=receiver_node.id if receiver_node else None, @@ -138,6 +152,17 @@ def _handle_message( ) session.add(message) + # Add first receiver to junction table + 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, + ) + 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 61ebbf4..d4e16ca 100644 --- a/src/meshcore_hub/collector/handlers/telemetry.py +++ b/src/meshcore_hub/collector/handlers/telemetry.py @@ -8,7 +8,7 @@ from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_telemetry_hash -from meshcore_hub.common.models import Node, Telemetry +from meshcore_hub.common.models import Node, Telemetry, add_event_receiver logger = logging.getLogger(__name__) @@ -58,18 +58,7 @@ def handle_telemetry( ) with db.session_scope() as session: - # Check if telemetry with same hash already exists - existing = session.execute( - select(Telemetry.id).where(Telemetry.event_hash == event_hash) - ).scalar_one_or_none() - - if existing: - logger.debug( - f"Duplicate telemetry skipped (node={node_public_key[:12]}...)" - ) - return - - # Find receiver node + # Find or create receiver node first (needed for both new and duplicate events) receiver_node = None if public_key: receiver_query = select(Node).where(Node.public_key == public_key) @@ -86,6 +75,29 @@ def handle_telemetry( else: receiver_node.last_seen = now + # Check if telemetry with same hash already exists + existing = session.execute( + select(Telemetry.id).where(Telemetry.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + # Event already exists - just add this receiver to the junction table + if receiver_node: + added = add_event_receiver( + session=session, + event_type="telemetry", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=None, + received_at=now, + ) + if added: + logger.debug( + f"Added receiver {public_key[:12]}... to telemetry " + f"(node={node_public_key[:12]}...)" + ) + return + # Find or create reporting node reporting_node = None if node_public_key: @@ -115,6 +127,17 @@ def handle_telemetry( ) session.add(telemetry) + # Add first receiver to junction table + 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, + ) + # 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 280f764..3d53d5b 100644 --- a/src/meshcore_hub/collector/handlers/trace.py +++ b/src/meshcore_hub/collector/handlers/trace.py @@ -8,7 +8,7 @@ from sqlalchemy import select from meshcore_hub.common.database import DatabaseManager from meshcore_hub.common.hash_utils import compute_trace_hash -from meshcore_hub.common.models import Node, TracePath +from meshcore_hub.common.models import Node, TracePath, add_event_receiver logger = logging.getLogger(__name__) @@ -45,16 +45,7 @@ def handle_trace_data( event_hash = compute_trace_hash(initiator_tag=initiator_tag) with db.session_scope() as session: - # Check if trace with same hash already exists - existing = session.execute( - select(TracePath.id).where(TracePath.event_hash == event_hash) - ).scalar_one_or_none() - - if existing: - logger.debug(f"Duplicate trace skipped (tag={initiator_tag})") - return - - # Find receiver node + # Find or create receiver node first (needed for both new and duplicate events) receiver_node = None if public_key: receiver_query = select(Node).where(Node.public_key == public_key) @@ -71,6 +62,29 @@ def handle_trace_data( else: receiver_node.last_seen = now + # Check if trace with same hash already exists + existing = session.execute( + select(TracePath.id).where(TracePath.event_hash == event_hash) + ).scalar_one_or_none() + + if existing: + # Event already exists - just add this receiver to the junction table + if receiver_node: + added = add_event_receiver( + session=session, + event_type="trace", + event_hash=event_hash, + receiver_node_id=receiver_node.id, + snr=None, # Trace events don't have a single SNR value + received_at=now, + ) + if added: + logger.debug( + f"Added receiver {public_key[:12]}... to trace " + f"(tag={initiator_tag})" + ) + return + # Create trace path record trace_path = TracePath( receiver_node_id=receiver_node.id if receiver_node else None, @@ -86,4 +100,15 @@ def handle_trace_data( ) session.add(trace_path) + # Add first receiver to junction table + 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, + ) + logger.info(f"Stored trace data: tag={initiator_tag}, hops={hop_count}") diff --git a/src/meshcore_hub/common/models/__init__.py b/src/meshcore_hub/common/models/__init__.py index d1dd8f0..69eaeb9 100644 --- a/src/meshcore_hub/common/models/__init__.py +++ b/src/meshcore_hub/common/models/__init__.py @@ -10,6 +10,7 @@ from meshcore_hub.common.models.telemetry import Telemetry from meshcore_hub.common.models.event_log import EventLog from meshcore_hub.common.models.member import Member from meshcore_hub.common.models.member_node import MemberNode +from meshcore_hub.common.models.event_receiver import EventReceiver, add_event_receiver __all__ = [ "Base", @@ -23,4 +24,6 @@ __all__ = [ "EventLog", "Member", "MemberNode", + "EventReceiver", + "add_event_receiver", ] diff --git a/src/meshcore_hub/common/models/event_receiver.py b/src/meshcore_hub/common/models/event_receiver.py new file mode 100644 index 0000000..c17f447 --- /dev/null +++ b/src/meshcore_hub/common/models/event_receiver.py @@ -0,0 +1,127 @@ +"""EventReceiver model for tracking which nodes received each event.""" + +from datetime import datetime +from typing import TYPE_CHECKING, Optional +from uuid import uuid4 + +from sqlalchemy import DateTime, Float, ForeignKey, Index, String, UniqueConstraint +from sqlalchemy.dialects.sqlite import insert as sqlite_insert +from sqlalchemy.orm import Mapped, Session, mapped_column, relationship + +from meshcore_hub.common.models.base import Base, TimestampMixin, UUIDMixin, utc_now + +if TYPE_CHECKING: + from meshcore_hub.common.models.node import Node + + +class EventReceiver(Base, UUIDMixin, TimestampMixin): + """Junction model tracking which receivers observed each event. + + This table enables multi-receiver tracking for deduplicated events. + When multiple receiver nodes observe the same mesh event, each receiver + gets an entry in this table linked by the event_hash. + + Attributes: + id: UUID primary key + event_type: Type of event ('message', 'advertisement', 'trace', 'telemetry') + event_hash: Hash identifying the unique event (links to event tables) + receiver_node_id: FK to the node that received this event + snr: Signal-to-noise ratio at this receiver (if available) + received_at: When this specific receiver saw the event + created_at: Record creation timestamp + updated_at: Record update timestamp + """ + + __tablename__ = "event_receivers" + + event_type: Mapped[str] = mapped_column( + String(20), + nullable=False, + ) + event_hash: Mapped[str] = mapped_column( + String(32), + nullable=False, + index=True, + ) + receiver_node_id: Mapped[str] = mapped_column( + String(36), + ForeignKey("nodes.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + snr: Mapped[Optional[float]] = mapped_column( + Float, + nullable=True, + ) + received_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), + default=utc_now, + nullable=False, + ) + + # Relationship to receiver node + receiver_node: Mapped["Node"] = relationship( + "Node", + foreign_keys=[receiver_node_id], + ) + + __table_args__ = ( + UniqueConstraint( + "event_hash", "receiver_node_id", name="uq_event_receivers_hash_node" + ), + Index("ix_event_receivers_type_hash", "event_type", "event_hash"), + ) + + def __repr__(self) -> str: + return ( + f"" + ) + + +def add_event_receiver( + session: Session, + event_type: str, + event_hash: str, + receiver_node_id: str, + snr: Optional[float] = None, + received_at: Optional[datetime] = None, +) -> bool: + """Add a receiver to an event, handling duplicates gracefully. + + Uses INSERT OR IGNORE to handle the unique constraint on (event_hash, receiver_node_id). + + Args: + session: SQLAlchemy session + event_type: Type of event ('message', 'advertisement', 'trace', 'telemetry') + event_hash: Hash identifying the unique event + receiver_node_id: UUID of the receiver node + snr: Signal-to-noise ratio at this receiver (optional) + received_at: When this receiver saw the event (defaults to now) + + Returns: + True if a new receiver entry was added, False if it already existed. + """ + from datetime import timezone + + now = received_at or datetime.now(timezone.utc) + + stmt = ( + sqlite_insert(EventReceiver) + .values( + id=str(uuid4()), + event_type=event_type, + event_hash=event_hash, + receiver_node_id=receiver_node_id, + snr=snr, + received_at=now, + created_at=now, + updated_at=now, + ) + .on_conflict_do_nothing(index_elements=["event_hash", "receiver_node_id"]) + ) + result = session.execute(stmt) + # CursorResult has rowcount attribute + rowcount = getattr(result, "rowcount", 0) + return bool(rowcount and rowcount > 0) diff --git a/src/meshcore_hub/common/schemas/__init__.py b/src/meshcore_hub/common/schemas/__init__.py index f9978e9..f9b2580 100644 --- a/src/meshcore_hub/common/schemas/__init__.py +++ b/src/meshcore_hub/common/schemas/__init__.py @@ -20,6 +20,7 @@ from meshcore_hub.common.schemas.nodes import ( NodeTagRead, ) from meshcore_hub.common.schemas.messages import ( + ReceiverInfo, MessageRead, MessageList, MessageFilters, @@ -57,7 +58,8 @@ __all__ = [ "NodeTagCreate", "NodeTagUpdate", "NodeTagRead", - # Messages + # Messages & Events + "ReceiverInfo", "MessageRead", "MessageList", "MessageFilters", diff --git a/src/meshcore_hub/common/schemas/messages.py b/src/meshcore_hub/common/schemas/messages.py index 6e74076..c19a876 100644 --- a/src/meshcore_hub/common/schemas/messages.py +++ b/src/meshcore_hub/common/schemas/messages.py @@ -6,6 +6,24 @@ from typing import Literal, Optional from pydantic import BaseModel, Field +class ReceiverInfo(BaseModel): + """Information about a receiver that observed an event.""" + + node_id: str = Field(..., description="Receiver node UUID") + public_key: str = Field(..., description="Receiver node public key") + name: Optional[str] = Field(default=None, description="Receiver node name") + friendly_name: Optional[str] = Field( + default=None, description="Receiver friendly name from tags" + ) + snr: Optional[float] = Field( + default=None, description="Signal-to-noise ratio at this receiver" + ) + received_at: datetime = Field(..., description="When this receiver saw the event") + + class Config: + from_attributes = True + + class MessageRead(BaseModel): """Schema for reading a message.""" @@ -37,6 +55,9 @@ class MessageRead(BaseModel): ) received_at: datetime = Field(..., description="When received by interface") created_at: datetime = Field(..., description="Record creation timestamp") + receivers: list[ReceiverInfo] = Field( + default_factory=list, description="All receivers that observed this message" + ) class Config: from_attributes = True @@ -104,6 +125,10 @@ class AdvertisementRead(BaseModel): flags: Optional[int] = Field(default=None, description="Capability flags") received_at: datetime = Field(..., description="When received") created_at: datetime = Field(..., description="Record creation timestamp") + receivers: list[ReceiverInfo] = Field( + default_factory=list, + description="All receivers that observed this advertisement", + ) class Config: from_attributes = True @@ -137,6 +162,10 @@ class TracePathRead(BaseModel): hop_count: Optional[int] = Field(default=None, description="Total hops") received_at: datetime = Field(..., description="When received") created_at: datetime = Field(..., description="Record creation timestamp") + receivers: list[ReceiverInfo] = Field( + default_factory=list, + description="All receivers that observed this trace", + ) class Config: from_attributes = True @@ -163,6 +192,10 @@ class TelemetryRead(BaseModel): ) received_at: datetime = Field(..., description="When received") created_at: datetime = Field(..., description="Record creation timestamp") + receivers: list[ReceiverInfo] = Field( + default_factory=list, + description="All receivers that observed this telemetry", + ) class Config: from_attributes = True diff --git a/src/meshcore_hub/web/templates/advertisements.html b/src/meshcore_hub/web/templates/advertisements.html index d3c5df8..41b6b06 100644 --- a/src/meshcore_hub/web/templates/advertisements.html +++ b/src/meshcore_hub/web/templates/advertisements.html @@ -71,7 +71,31 @@ {% endif %} - {% if ad.received_by %} + {% if ad.receivers and ad.receivers|length > 1 %} + + {% elif ad.receivers and ad.receivers|length == 1 %} + + {% if ad.receivers[0].friendly_name or ad.receivers[0].name %} +
{{ ad.receivers[0].friendly_name or ad.receivers[0].name }}
+
{{ ad.receivers[0].public_key[:16] }}...
+ {% else %} + {{ ad.receivers[0].public_key[:16] }}... + {% endif %} +
+ {% elif ad.received_by %} {% if ad.receiver_friendly_name or ad.receiver_name %}
{{ ad.receiver_friendly_name or ad.receiver_name }}
diff --git a/src/meshcore_hub/web/templates/messages.html b/src/meshcore_hub/web/templates/messages.html index 8fe1425..2c68e29 100644 --- a/src/meshcore_hub/web/templates/messages.html +++ b/src/meshcore_hub/web/templates/messages.html @@ -87,7 +87,34 @@ {{ msg.text or '-' }} - {% if msg.received_by %} + {% if msg.receivers and msg.receivers|length > 1 %} +
+ {% elif msg.receivers and msg.receivers|length == 1 %} + + {% if msg.receivers[0].friendly_name or msg.receivers[0].name %} +
{{ msg.receivers[0].friendly_name or msg.receivers[0].name }}
+
{{ msg.receivers[0].public_key[:16] }}...
+ {% else %} + {{ msg.receivers[0].public_key[:16] }}... + {% endif %} +
+ {% elif msg.received_by %} {% if msg.receiver_friendly_name or msg.receiver_name %}
{{ msg.receiver_friendly_name or msg.receiver_name }}
diff --git a/tests/test_common/test_hash_utils.py b/tests/test_common/test_hash_utils.py index 1ca94c4..62a85c2 100644 --- a/tests/test_common/test_hash_utils.py +++ b/tests/test_common/test_hash_utils.py @@ -94,7 +94,7 @@ class TestComputeAdvertisementHash: def test_same_content_same_bucket_produces_same_hash(self) -> None: """Advertisements within the same time bucket should match.""" - # Two times within the same 5-minute bucket + # Two times within the same 5-minute (300 second) bucket time1 = datetime(2024, 1, 15, 10, 31, 0, tzinfo=timezone.utc) time2 = datetime(2024, 1, 15, 10, 33, 0, tzinfo=timezone.utc) @@ -104,7 +104,7 @@ class TestComputeAdvertisementHash: adv_type="chat", flags=128, received_at=time1, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) hash2 = compute_advertisement_hash( public_key="a" * 64, @@ -112,14 +112,14 @@ class TestComputeAdvertisementHash: adv_type="chat", flags=128, received_at=time2, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) assert hash1 == hash2 def test_different_bucket_produces_different_hash(self) -> None: """Advertisements in different time buckets should not match.""" - # Two times in different 5-minute buckets + # Two times in different 5-minute (300 second) buckets time1 = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) time2 = datetime(2024, 1, 15, 10, 36, 0, tzinfo=timezone.utc) @@ -127,13 +127,13 @@ class TestComputeAdvertisementHash: public_key="a" * 64, name="Node1", received_at=time1, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) hash2 = compute_advertisement_hash( public_key="a" * 64, name="Node1", received_at=time2, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) assert hash1 != hash2 @@ -158,29 +158,29 @@ class TestComputeAdvertisementHash: time1 = datetime(2024, 1, 15, 10, 30, 0, tzinfo=timezone.utc) time2 = datetime(2024, 1, 15, 10, 35, 0, tzinfo=timezone.utc) - # With 5-minute bucket, these should be in different buckets + # With 5-minute (300s) bucket, these should be in different buckets hash1_5min = compute_advertisement_hash( public_key="a" * 64, received_at=time1, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) hash2_5min = compute_advertisement_hash( public_key="a" * 64, received_at=time2, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) assert hash1_5min != hash2_5min - # With 10-minute bucket, these should be in the same bucket + # With 10-minute (600s) bucket, these should be in the same bucket hash1_10min = compute_advertisement_hash( public_key="a" * 64, received_at=time1, - bucket_minutes=10, + bucket_seconds=600, # 10 minutes ) hash2_10min = compute_advertisement_hash( public_key="a" * 64, received_at=time2, - bucket_minutes=10, + bucket_seconds=600, # 10 minutes ) assert hash1_10min == hash2_10min @@ -216,13 +216,13 @@ class TestComputeTelemetryHash: node_public_key="a" * 64, parsed_data=data, received_at=time1, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) hash2 = compute_telemetry_hash( node_public_key="a" * 64, parsed_data=data, received_at=time2, - bucket_minutes=5, + bucket_seconds=300, # 5 minutes ) assert hash1 == hash2 From 733342a9eccc719e5a25660fa8965a3a724cde85 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 14:21:17 +0000 Subject: [PATCH 05/10] Fixed README and Compose --- README.md | 44 +++++++++++++++++--------------------------- docker-compose.yml | 16 ++++++++++++++-- 2 files changed, 31 insertions(+), 29 deletions(-) diff --git a/README.md b/README.md index f8e4d29..51c4929 100644 --- a/README.md +++ b/README.md @@ -83,7 +83,7 @@ cp .env.example .env # Edit .env: set SERIAL_PORT to your device (e.g., /dev/ttyUSB0 or /dev/ttyACM0) # Start the entire stack including the interface receiver -docker compose --profile interface-receiver up -d +docker compose --profile core --profile receiver up -d # View the web dashboard open http://localhost:8080 @@ -105,7 +105,7 @@ For larger deployments, you can separate receiver nodes from the central infrast │ │ Device │ │ Device │ │ Device │ │ │ └──────┬───────┘ └──────┬───────┘ └──────┬───────┘ │ │ │ │ │ │ -│ │ interface-receiver only │ │ +│ │ receiver profile only │ │ │ └──────────────────┼──────────────────┘ │ │ │ │ │ MQTT (port 1883) │ @@ -126,18 +126,18 @@ For larger deployments, you can separate receiver nodes from the central infrast **On each receiver node (Raspberry Pi, etc.):** ```bash -# Only run the interface-receiver component +# Only run the receiver component # Configure .env with MQTT_HOST pointing to your central server MQTT_HOST=your-community-server.com SERIAL_PORT=/dev/ttyUSB0 -docker compose --profile interface-receiver up -d +docker compose --profile receiver up -d ``` **On the central server (VPS/cloud):** ```bash # Run the core infrastructure (no interface needed) -docker compose up -d +docker compose --profile core up -d ``` This architecture allows: @@ -150,26 +150,16 @@ This architecture allows: ### Using Docker Compose (Recommended) -Docker Compose runs core services by default and uses **profiles** for optional components: +Docker Compose uses **profiles** to select which services to run: -**Default Services (always run):** - -| Service | Description | -|---------|-------------| -| `mqtt` | Eclipse Mosquitto MQTT broker | -| `collector` | MQTT subscriber + database storage (auto-seeds on startup) | -| `api` | REST API server | -| `web` | Web dashboard | - -**Optional Profiles:** - -| Profile | Services | -|---------|----------| -| `interface-receiver` | MeshCore device receiver (events to MQTT) | -| `interface-sender` | MeshCore device sender (MQTT to device) | -| `mock` | Mock device receiver (for testing without hardware) | -| `migrate` | One-time database migration runner | -| `seed` | One-time seed data import (also runs automatically on collector startup) | +| Profile | Services | Use Case | +|---------|----------|----------| +| `core` | mqtt, collector, api, web | Central server infrastructure | +| `receiver` | mqtt, interface-receiver | Receiver node (events to MQTT) | +| `sender` | mqtt, interface-sender | Sender node (MQTT to device) | +| `mock` | mqtt, interface-mock-receiver | Testing without hardware | +| `migrate` | db-migrate | One-time database migration | +| `seed` | seed | One-time seed data import | ```bash # Clone the repository @@ -187,10 +177,10 @@ docker compose --profile migrate run --rm db-migrate docker compose --profile seed run --rm seed # Start core services (mqtt, collector, api, web) -docker compose up -d +docker compose --profile core up -d -# Start sender/receiver interface -docker compose --profile interface-receiver up -d +# Or start just the receiver (for distributed setups) +docker compose --profile receiver up -d # View logs docker compose logs -f diff --git a/docker-compose.yml b/docker-compose.yml index 7a2a2a8..b012e95 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,10 +1,16 @@ services: # ========================================================================== # MQTT Broker - Eclipse Mosquitto + # Included in all profiles that need it (receiver, sender, core, mock) # ========================================================================== mqtt: image: eclipse-mosquitto:2 container_name: meshcore-mqtt + profiles: + - receiver + - sender + - core + - mock restart: unless-stopped ports: - "${MQTT_EXTERNAL_PORT:-1883}:1883" @@ -30,7 +36,7 @@ services: dockerfile: Dockerfile container_name: meshcore-interface-receiver profiles: - - interface-receiver + - receiver restart: unless-stopped depends_on: mqtt: @@ -66,7 +72,7 @@ services: dockerfile: Dockerfile container_name: meshcore-interface-sender profiles: - - interface-sender + - sender restart: unless-stopped depends_on: mqtt: @@ -133,6 +139,8 @@ services: context: . dockerfile: Dockerfile container_name: meshcore-collector + profiles: + - core restart: unless-stopped depends_on: mqtt: @@ -180,6 +188,8 @@ services: context: . dockerfile: Dockerfile container_name: meshcore-api + profiles: + - core restart: unless-stopped depends_on: mqtt: @@ -222,6 +232,8 @@ services: context: . dockerfile: Dockerfile container_name: meshcore-web + profiles: + - core restart: unless-stopped depends_on: api: From 9cd1d50bf6b3f4644d1ad9193caf40d9b96903b4 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 14:38:53 +0000 Subject: [PATCH 06/10] Updates --- README.md | 27 ++++++++++++++++++--------- docker-compose.yml | 23 +++-------------------- 2 files changed, 21 insertions(+), 29 deletions(-) diff --git a/README.md b/README.md index 51c4929..6c230df 100644 --- a/README.md +++ b/README.md @@ -82,8 +82,8 @@ cd meshcore-hub cp .env.example .env # Edit .env: set SERIAL_PORT to your device (e.g., /dev/ttyUSB0 or /dev/ttyACM0) -# Start the entire stack including the interface receiver -docker compose --profile core --profile receiver up -d +# Start the entire stack with local MQTT broker +docker compose --profile mqtt --profile core --profile receiver up -d # View the web dashboard open http://localhost:8080 @@ -136,7 +136,10 @@ docker compose --profile receiver up -d **On the central server (VPS/cloud):** ```bash -# Run the core infrastructure (no interface needed) +# Run the core infrastructure with local MQTT broker +docker compose --profile mqtt --profile core up -d + +# Or connect to an existing MQTT broker (set MQTT_HOST in .env) docker compose --profile core up -d ``` @@ -154,13 +157,16 @@ Docker Compose uses **profiles** to select which services to run: | Profile | Services | Use Case | |---------|----------|----------| -| `core` | mqtt, collector, api, web | Central server infrastructure | -| `receiver` | mqtt, interface-receiver | Receiver node (events to MQTT) | -| `sender` | mqtt, interface-sender | Sender node (MQTT to device) | -| `mock` | mqtt, interface-mock-receiver | Testing without hardware | +| `core` | collector, api, web | Central server infrastructure | +| `receiver` | interface-receiver | Receiver node (events to MQTT) | +| `sender` | interface-sender | Sender node (MQTT to device) | +| `mqtt` | mosquitto broker | Local MQTT broker (optional) | +| `mock` | interface-mock-receiver | Testing without hardware | | `migrate` | db-migrate | One-time database migration | | `seed` | seed | One-time seed data import | +**Note:** Most deployments connect to an external MQTT broker. Add `--profile mqtt` only if you need a local broker. + ```bash # Clone the repository git clone https://github.com/ipnet-mesh/meshcore-hub.git @@ -176,10 +182,13 @@ docker compose --profile migrate run --rm db-migrate # Seed the database docker compose --profile seed run --rm seed -# Start core services (mqtt, collector, api, web) +# Start core services with local MQTT broker +docker compose --profile mqtt --profile core up -d + +# Or connect to external MQTT (configure MQTT_HOST in .env) docker compose --profile core up -d -# Or start just the receiver (for distributed setups) +# Start just the receiver (connects to MQTT_HOST from .env) docker compose --profile receiver up -d # View logs diff --git a/docker-compose.yml b/docker-compose.yml index b012e95..80e715c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,16 +1,13 @@ services: # ========================================================================== - # MQTT Broker - Eclipse Mosquitto - # Included in all profiles that need it (receiver, sender, core, mock) + # MQTT Broker - Eclipse Mosquitto (optional, use --profile mqtt) + # Most users will connect to an external MQTT broker instead # ========================================================================== mqtt: image: eclipse-mosquitto:2 container_name: meshcore-mqtt profiles: - - receiver - - sender - - core - - mock + - mqtt restart: unless-stopped ports: - "${MQTT_EXTERNAL_PORT:-1883}:1883" @@ -38,9 +35,6 @@ services: profiles: - receiver restart: unless-stopped - depends_on: - mqtt: - condition: service_healthy devices: - "${SERIAL_PORT:-/dev/ttyUSB0}:${SERIAL_PORT:-/dev/ttyUSB0}" user: root # Required for device access @@ -74,9 +68,6 @@ services: profiles: - sender restart: unless-stopped - depends_on: - mqtt: - condition: service_healthy devices: - "${SERIAL_PORT_SENDER:-/dev/ttyUSB1}:${SERIAL_PORT_SENDER:-/dev/ttyUSB1}" user: root # Required for device access @@ -110,9 +101,6 @@ services: profiles: - mock restart: unless-stopped - depends_on: - mqtt: - condition: service_healthy environment: - LOG_LEVEL=${LOG_LEVEL:-INFO} - MQTT_HOST=${MQTT_HOST:-mqtt} @@ -142,9 +130,6 @@ services: profiles: - core restart: unless-stopped - depends_on: - mqtt: - condition: service_healthy volumes: - ${DATA_HOME:-./data}:/data - ${SEED_HOME:-./seed}:/seed @@ -192,8 +177,6 @@ services: - core restart: unless-stopped depends_on: - mqtt: - condition: service_healthy collector: condition: service_started ports: From 78a086e1ea1e965be5bbbe6b61ed316d060f0b37 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 14:51:47 +0000 Subject: [PATCH 07/10] 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: From dd0b4c73c5fd85cc175354c27e4b49692481f1ce Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 15:10:03 +0000 Subject: [PATCH 08/10] More fixes --- docker-compose.yml | 9 ++++++++ src/meshcore_hub/interface/receiver.py | 31 ++++++++++++++++---------- src/meshcore_hub/interface/sender.py | 25 +++++++++++++-------- 3 files changed, 44 insertions(+), 21 deletions(-) diff --git a/docker-compose.yml b/docker-compose.yml index 80e715c..a42c58a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -7,6 +7,7 @@ services: image: eclipse-mosquitto:2 container_name: meshcore-mqtt profiles: + - all - mqtt restart: unless-stopped ports: @@ -33,6 +34,7 @@ services: dockerfile: Dockerfile container_name: meshcore-interface-receiver profiles: + - all - receiver restart: unless-stopped devices: @@ -66,6 +68,7 @@ services: dockerfile: Dockerfile container_name: meshcore-interface-sender profiles: + - all - sender restart: unless-stopped devices: @@ -99,6 +102,7 @@ services: dockerfile: Dockerfile container_name: meshcore-interface-mock-receiver profiles: + - all - mock restart: unless-stopped environment: @@ -128,6 +132,7 @@ services: dockerfile: Dockerfile container_name: meshcore-collector profiles: + - all - core restart: unless-stopped volumes: @@ -174,6 +179,7 @@ services: dockerfile: Dockerfile container_name: meshcore-api profiles: + - all - core restart: unless-stopped depends_on: @@ -216,6 +222,7 @@ services: dockerfile: Dockerfile container_name: meshcore-web profiles: + - all - core restart: unless-stopped depends_on: @@ -255,6 +262,7 @@ services: dockerfile: Dockerfile container_name: meshcore-db-migrate profiles: + - all - migrate volumes: # Mount data directory (uses collector/meshcore.db) @@ -275,6 +283,7 @@ services: dockerfile: Dockerfile container_name: meshcore-seed profiles: + - all - seed volumes: # Mount data directory for database (read-write) diff --git a/src/meshcore_hub/interface/receiver.py b/src/meshcore_hub/interface/receiver.py index ec57805..7a86044 100644 --- a/src/meshcore_hub/interface/receiver.py +++ b/src/meshcore_hub/interface/receiver.py @@ -197,17 +197,19 @@ class Receiver: logger.error(f"Failed to connect to MQTT broker: {e}") raise - # Connect to device - if not self.device.connect(): - self._device_connected = False - logger.error("Failed to connect to MeshCore device") - self.mqtt.stop() - self.mqtt.disconnect() - self._mqtt_connected = False - raise RuntimeError("Failed to connect to MeshCore device") + # Device should already be connected (from create_receiver) + # but handle case where start() is called directly + if not self.device.is_connected: + if not self.device.connect(): + self._device_connected = False + logger.error("Failed to connect to MeshCore device") + self.mqtt.stop() + self.mqtt.disconnect() + self._mqtt_connected = False + raise RuntimeError("Failed to connect to MeshCore device") + logger.info(f"Connected to MeshCore device: {self.device.public_key}") self._device_connected = True - logger.info(f"Connected to MeshCore device: {self.device.public_key}") # Initialize device: set time and send local advertisement self._initialize_device() @@ -291,17 +293,22 @@ def create_receiver( Returns: Configured Receiver instance """ - # Create device + # Create and connect device first to get public key device = create_device(port=port, baud=baud, mock=mock, node_address=node_address) - # Create MQTT client + if not device.connect(): + raise RuntimeError("Failed to connect to MeshCore device") + + logger.info(f"Connected to MeshCore device: {device.public_key}") + + # Create MQTT client with device's public key for unique client ID mqtt_config = MQTTConfig( host=mqtt_host, port=mqtt_port, username=mqtt_username, password=mqtt_password, prefix=mqtt_prefix, - client_id=f"meshcore-receiver-{device.public_key[:8] if device.public_key else 'unknown'}", + client_id=f"meshcore-receiver-{device.public_key[:12] if device.public_key else 'unknown'}", ) mqtt_client = MQTTClient(mqtt_config) diff --git a/src/meshcore_hub/interface/sender.py b/src/meshcore_hub/interface/sender.py index 75f280c..4415a32 100644 --- a/src/meshcore_hub/interface/sender.py +++ b/src/meshcore_hub/interface/sender.py @@ -200,14 +200,16 @@ class Sender: """Start the sender.""" logger.info("Starting SENDER mode") - # Connect to device first - if not self.device.connect(): - self._device_connected = False - logger.error("Failed to connect to MeshCore device") - raise RuntimeError("Failed to connect to MeshCore device") + # Device should already be connected (from create_sender) + # but handle case where start() is called directly + if not self.device.is_connected: + if not self.device.connect(): + self._device_connected = False + logger.error("Failed to connect to MeshCore device") + raise RuntimeError("Failed to connect to MeshCore device") + logger.info(f"Connected to MeshCore device: {self.device.public_key}") self._device_connected = True - logger.info(f"Connected to MeshCore device: {self.device.public_key}") # Connect to MQTT broker try: @@ -307,17 +309,22 @@ def create_sender( Returns: Configured Sender instance """ - # Create device + # Create and connect device first to get public key device = create_device(port=port, baud=baud, mock=mock, node_address=node_address) - # Create MQTT client + if not device.connect(): + raise RuntimeError("Failed to connect to MeshCore device") + + logger.info(f"Connected to MeshCore device: {device.public_key}") + + # Create MQTT client with device's public key for unique client ID mqtt_config = MQTTConfig( host=mqtt_host, port=mqtt_port, username=mqtt_username, password=mqtt_password, prefix=mqtt_prefix, - client_id=f"meshcore-sender-{device.public_key[:8] if device.public_key else 'unknown'}", + client_id=f"meshcore-sender-{device.public_key[:12] if device.public_key else 'unknown'}", ) mqtt_client = MQTTClient(mqtt_config) From f0531c9e407c98ebcb017cf578d18e2c1f2cc200 Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 15:16:26 +0000 Subject: [PATCH 09/10] Updated env example --- .env.example | 51 ++++++++++++++++++++++++++------------------------- 1 file changed, 26 insertions(+), 25 deletions(-) diff --git a/.env.example b/.env.example index 2248e07..e3d4e36 100644 --- a/.env.example +++ b/.env.example @@ -2,32 +2,28 @@ # Copy this file to .env and customize values # =================== -# Docker Image +# Data & Seed Directories # =================== -# Leave empty to build from local Dockerfile, or set to use a pre-built image: -# MESHCORE_IMAGE=ghcr.io/ipnet-mesh/meshcore-hub:latest -# MESHCORE_IMAGE=ghcr.io/ipnet-mesh/meshcore-hub:main -# MESHCORE_IMAGE=ghcr.io/ipnet-mesh/meshcore-hub:v1.0.0 -MESHCORE_IMAGE= - -# =================== -# Data Directory -# =================== - -# Base directory for all service data (collector DB, tags, members, etc.) +# Base directory for runtime data (database, etc.) # Default: ./data (relative to docker-compose.yml location) # Inside containers this is mapped to /data # # Structure: # ${DATA_HOME}/ -# ├── collector/ -# │ ├── meshcore.db # SQLite database -# │ └── tags.json # Node tags for import -# └── web/ -# └── members.json # Network members list +# └── meshcore.db # SQLite database DATA_HOME=./data +# Directory containing seed data files for import +# Default: ./seed (relative to docker-compose.yml location) +# Inside containers this is mapped to /seed +# +# Structure: +# ${SEED_HOME}/ +# ├── node_tags.yaml # Node tags for import +# └── members.yaml # Network members for import +SEED_HOME=./seed + # =================== # Common Settings # =================== @@ -35,12 +31,20 @@ DATA_HOME=./data # Logging level (DEBUG, INFO, WARNING, ERROR, CRITICAL) LOG_LEVEL=INFO -# MQTT Broker Settings (internal use) +# =================== +# MQTT Settings +# =================== + +# MQTT Broker connection (for interface/collector/api services) +# When using the local MQTT broker (--profile mqtt), use "mqtt" as host +# When using an external broker, set the hostname/IP +MQTT_HOST=mqtt +MQTT_PORT=1883 MQTT_USERNAME= MQTT_PASSWORD= MQTT_PREFIX=meshcore -# External MQTT port mapping +# External port mappings for local MQTT broker (--profile mqtt only) MQTT_EXTERNAL_PORT=1883 MQTT_WS_PORT=9001 @@ -58,6 +62,7 @@ SERIAL_PORT_SENDER=/dev/ttyUSB1 SERIAL_BAUD=115200 # Optional node address override (64-char hex string) +# Only set if you need to override the device's public key NODE_ADDRESS= NODE_ADDRESS_SENDER= @@ -84,25 +89,21 @@ WEB_PORT=8080 NETWORK_NAME=MeshCore Network NETWORK_CITY= NETWORK_COUNTRY= -NETWORK_LOCATION= # Radio configuration (comma-delimited) # Format: ,,,,, # Example: EU/UK Narrow,869.618MHz,62.5kHz,8,8,22dBm NETWORK_RADIO_CONFIG= +# Contact information NETWORK_CONTACT_EMAIL= NETWORK_CONTACT_DISCORD= +NETWORK_CONTACT_GITHUB= # Welcome text displayed on the homepage (plain text, optional) # If not set, a default welcome message is shown NETWORK_WELCOME_TEXT= -# Members file location (optional override) -# Default: ${DATA_HOME}/web/members.json -# Only set this if you want to use a different location -# MEMBERS_FILE=/custom/path/to/members.json - # =================== # Webhook Settings # =================== From 714c3cbbd274f81478a03a1a4fa041c410a0773d Mon Sep 17 00:00:00 2001 From: Louis King Date: Sat, 6 Dec 2025 15:32:15 +0000 Subject: [PATCH 10/10] Set sensible Docker tag label --- .env.example | 8 ++++++++ docker-compose.yml | 16 ++++++++-------- 2 files changed, 16 insertions(+), 8 deletions(-) diff --git a/.env.example b/.env.example index e3d4e36..f3428b1 100644 --- a/.env.example +++ b/.env.example @@ -1,6 +1,14 @@ # MeshCore Hub - Docker Compose Environment Configuration # Copy this file to .env and customize values +# =================== +# Docker Image +# =================== + +# Docker image version tag to use +# Options: latest, main, v1.0.0, etc. +IMAGE_VERSION=latest + # =================== # Data & Seed Directories # =================== diff --git a/docker-compose.yml b/docker-compose.yml index a42c58a..cb074b0 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -28,7 +28,7 @@ services: # Interface Receiver - MeshCore device to MQTT bridge (events) # ========================================================================== interface-receiver: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -62,7 +62,7 @@ services: # Interface Sender - MQTT to MeshCore device bridge (commands) # ========================================================================== interface-sender: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -96,7 +96,7 @@ services: # Interface Mock Receiver - For testing without real devices # ========================================================================== interface-mock-receiver: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -126,7 +126,7 @@ services: # Collector - MQTT subscriber and database storage # ========================================================================== collector: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -173,7 +173,7 @@ services: # API Server - REST API for querying data and sending commands # ========================================================================== api: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -216,7 +216,7 @@ services: # Web Dashboard - Web interface for network visualization # ========================================================================== web: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -256,7 +256,7 @@ services: # Database Migrations - Run Alembic migrations # ========================================================================== db-migrate: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile @@ -277,7 +277,7 @@ services: # Seed Data - Import node_tags.json and members.json from SEED_HOME # ========================================================================== seed: - image: ghcr.io/ipnet-mesh/meshcore-hub:main + image: ghcr.io/ipnet-mesh/meshcore-hub:${IMAGE_VERSION:-latest} build: context: . dockerfile: Dockerfile