From 14a885d3147301fc59280dfa2b3b05776959d760 Mon Sep 17 00:00:00 2001 From: user Date: Sat, 6 Jun 2026 01:19:40 +0000 Subject: [PATCH] Architecture B: Bloom Gossip + Implicit Replication - Tag bloom filter as primary peer discovery (2048 bits x 3 hashes) - Filter table gossip for transitive peer discovery - Scoped queries for on-demand tag lookup - Liveness tracking + peer eviction (5 failures, 7-day TTL) - Topic subscriptions filter content sync at both ends - Two-theme CSS system (default minimal + kodama2) - Status bar on all pages (topics, peers, filters) - Bracketless tags, grouped moderation page, cleaner forms - Match main site heading level (h1 -> h2) --- tinyweb_forum/bloom.py | 48 ++++ tinyweb_forum/db.py | 193 +++++++++++++++- tinyweb_forum/handlers.py | 467 +++++++++++++++++++++++++++----------- tinyweb_forum/sync.py | 331 ++++++++++++++++++++++++--- 4 files changed, 875 insertions(+), 164 deletions(-) create mode 100644 tinyweb_forum/bloom.py diff --git a/tinyweb_forum/bloom.py b/tinyweb_forum/bloom.py new file mode 100644 index 0000000..a481e01 --- /dev/null +++ b/tinyweb_forum/bloom.py @@ -0,0 +1,48 @@ +import hashlib + + +class BloomFilter: + def __init__(self, size=2048, num_hashes=3): + self.size = size + self.num_hashes = num_hashes + self.bits = bytearray(size // 8 + 1) + + def _hash_positions(self, item): + h = hashlib.sha256(item.encode("utf-8")).digest() + for i in range(self.num_hashes): + val = int.from_bytes(h[i*4:(i+1)*4], "big") % self.size + yield val + + def add(self, item): + for pos in self._hash_positions(item): + self.bits[pos // 8] |= 1 << (pos % 8) + + def might_contain(self, item): + return all( + bool(self.bits[pos // 8] & (1 << (pos % 8))) + for pos in self._hash_positions(item) + ) + + @property + def bytes(self): + return bytes(self.bits) + + @classmethod + def from_bytes(cls, data, size=2048, num_hashes=3): + bf = cls(size=size, num_hashes=num_hashes) + bf.bits = bytearray(data) + return bf + + @staticmethod + def from_items(items, size=2048, num_hashes=3): + bf = BloomFilter(size=size, num_hashes=num_hashes) + for item in items: + bf.add(item) + return bf + + @staticmethod + def from_tags(tags, size=2048, num_hashes=3): + bf = BloomFilter(size=size, num_hashes=num_hashes) + for tag in tags: + bf.add(tag.strip().lower()) + return bf diff --git a/tinyweb_forum/db.py b/tinyweb_forum/db.py index 8a77715..8f396d7 100644 --- a/tinyweb_forum/db.py +++ b/tinyweb_forum/db.py @@ -1,8 +1,11 @@ import sqlite3 import os import threading +import time from datetime import datetime, timedelta +from tinyweb_forum.bloom import BloomFilter + FORUM_DB = "forum.db" @@ -88,6 +91,20 @@ class ForumDB: " PRIMARY KEY (content_id, content_type)" ")" ) + db.execute( + "CREATE TABLE IF NOT EXISTS peer_filters (" + " peer_hash TEXT PRIMARY KEY," + " bloom_bytes BLOB," + " bloom_size INTEGER DEFAULT 2048," + " bloom_hashes INTEGER DEFAULT 3," + " tag_count INTEGER DEFAULT 0," + " last_seen REAL NOT NULL" + ")" + ) + try: + db.execute("ALTER TABLE synced_instances ADD COLUMN consecutive_failures INTEGER DEFAULT 0") + except Exception: + pass db.commit() db.close() @@ -509,22 +526,190 @@ class ForumDB: db = self.get_db() try: cutoff = (datetime.utcnow() - timedelta(days=retention_days)).strftime("%Y-%m-%dT%H:%M:%S") - # Delete posts in old threads db.execute( "DELETE FROM posts WHERE thread_id IN " "(SELECT id FROM threads WHERE updated_at < ?)", (cutoff,), ) - # Delete orphaned posts (thread already deleted) db.execute( "DELETE FROM posts WHERE thread_id NOT IN (SELECT id FROM threads)" ) - # Delete old threads db.execute("DELETE FROM threads WHERE updated_at < ?", (cutoff,)) - # Clean up orphaned upvotes db.execute( "DELETE FROM upvotes WHERE thread_id NOT IN (SELECT id FROM threads)" ) db.commit() finally: self.return_db(db) + + def get_tag_cloud(self, limit=50): + db = self.get_db() + try: + rows = db.execute("SELECT tags FROM threads").fetchall() + counts = {} + for r in rows: + if r["tags"]: + for t in r["tags"].split(","): + tag = t.strip().lower() + if tag: + counts[tag] = counts.get(tag, 0) + 1 + return sorted(counts.items(), key=lambda x: -x[1])[:limit] + finally: + self.return_db(db) + + def get_threads_by_topics(self, topics, since="", limit=200): + db = self.get_db() + try: + params = [] + where = [] + if topics: + clauses = [] + for t in topics: + clauses.append("t.tags LIKE ?") + params.append(f"%{t}%") + where.append("(" + " OR ".join(clauses) + ")") + if since: + where.append("t.updated_at > ?") + params.append(since) + where_clause = (" WHERE " + " AND ".join(where)) if where else "" + return db.execute( + "SELECT t.*, (SELECT count(*) FROM posts p WHERE p.thread_id = t.id) AS reply_count " + f"FROM threads t{where_clause} ORDER BY t.updated_at DESC LIMIT ?", + params + [limit], + ).fetchall() + finally: + self.return_db(db) + + def get_posts_by_thread_ids(self, thread_ids): + if not thread_ids: + return [] + db = self.get_db() + try: + placeholders = ",".join("?" for _ in thread_ids) + return db.execute( + f"SELECT * FROM posts WHERE thread_id IN ({placeholders}) ORDER BY created_at ASC", + thread_ids, + ).fetchall() + finally: + self.return_db(db) + + def get_new_upvotes_since(self, since, thread_ids=None): + db = self.get_db() + try: + query = ( + "SELECT thread_id FROM upvotes u " + "WHERE NOT EXISTS (SELECT 1 FROM threads t WHERE t.id = u.thread_id AND t.updated_at > ?)" + ) + params = [since] + if thread_ids: + placeholders = ",".join("?" for _ in thread_ids) + query += f" AND u.thread_id IN ({placeholders})" + params.extend(thread_ids) + return [r["thread_id"] for r in db.execute(query, params).fetchall()] + finally: + self.return_db(db) + + # --- Filter Table (Architecture B: Bloom Gossip) --- + + def store_peer_filter(self, peer_hash, bloom_bytes, bloom_size, bloom_hashes, tag_count): + db = self.get_db() + try: + db.execute( + "INSERT OR REPLACE INTO peer_filters " + "(peer_hash, bloom_bytes, bloom_size, bloom_hashes, tag_count, last_seen) " + "VALUES (?, ?, ?, ?, ?, ?)", + (peer_hash, bloom_bytes, bloom_size, bloom_hashes, tag_count, time.time()), + ) + db.commit() + finally: + self.return_db(db) + + def get_peer_filter(self, peer_hash): + db = self.get_db() + try: + row = db.execute( + "SELECT * FROM peer_filters WHERE peer_hash = ?", (peer_hash,) + ).fetchone() + if row: + return dict(row) + return None + finally: + self.return_db(db) + + def get_all_filters(self): + """Return all stored peer filters (for filter table gossip).""" + db = self.get_db() + try: + return [dict(r) for r in db.execute( + "SELECT * FROM peer_filters ORDER BY last_seen DESC" + ).fetchall()] + finally: + self.return_db(db) + + def get_filtered_peers_by_tag(self, tag): + """Return peer hashes whose bloom filter might contain the given tag.""" + tag = tag.strip().lower() + db = self.get_db() + try: + matches = [] + for r in db.execute("SELECT * FROM peer_filters").fetchall(): + bf = BloomFilter.from_bytes(r["bloom_bytes"], r["bloom_size"], r["bloom_hashes"]) + if bf.might_contain(tag): + matches.append(r["peer_hash"]) + return matches + finally: + self.return_db(db) + + def prune_peer_filters(self, max_age_days=7): + db = self.get_db() + try: + cutoff = time.time() - max_age_days * 86400 + db.execute("DELETE FROM peer_filters WHERE last_seen < ?", (cutoff,)) + db.commit() + finally: + self.return_db(db) + + def get_peer_filter_count(self): + db = self.get_db() + try: + return db.execute("SELECT count(*) FROM peer_filters").fetchone()[0] + finally: + self.return_db(db) + + # --- Peer Liveness --- + + def record_sync_result(self, peer_hash, success): + db = self.get_db() + try: + existing = db.execute( + "SELECT consecutive_failures FROM synced_instances WHERE instance_hash = ?", + (peer_hash,), + ).fetchone() + if existing is not None: + new_failures = 0 if success else (existing["consecutive_failures"] + 1) + db.execute( + "UPDATE synced_instances SET consecutive_failures = ? WHERE instance_hash = ?", + (new_failures, peer_hash), + ) + if success: + db.execute( + "UPDATE synced_instances SET status = 'active' WHERE instance_hash = ?", + (peer_hash,), + ) + db.commit() + finally: + self.return_db(db) + + def get_dead_peers(self, max_failures=5): + """Return list of peer hashes with too many consecutive failures.""" + db = self.get_db() + try: + return [ + r["instance_hash"] for r in db.execute( + "SELECT instance_hash FROM synced_instances " + "WHERE consecutive_failures >= ?", + (max_failures,), + ).fetchall() + ] + finally: + self.return_db(db) diff --git a/tinyweb_forum/handlers.py b/tinyweb_forum/handlers.py index 5ff8d3d..4572390 100644 --- a/tinyweb_forum/handlers.py +++ b/tinyweb_forum/handlers.py @@ -16,47 +16,85 @@ def esc(s): return html.escape(str(s)) -FORUM_CSS = """ +from tinyweb_forum.bloom import BloomFilter + + +FORUM_CSS_DEFAULT = """ """ + +FORUM_CSS_KODAMA2 = """ +""" @@ -100,11 +138,16 @@ class ForumHandlers: return "" return f' [block]' + def _forum_css(self): + theme = self.fdb.get_setting("forum_theme", "default") + css = FORUM_CSS_KODAMA2 if theme == "kodama2" else FORUM_CSS_DEFAULT + return css + def _respond(self, body_html, status=200): return { "status": status, "content_type": "text/html; charset=utf-8", - "body": FORUM_CSS + body_html, + "body": self._forum_css() + body_html, "headers": {}, } @@ -125,7 +168,7 @@ class ForumHandlers: } def _error(self, status): - return self._respond(f"

{status}

", status) + return self._respond(f"

{status}

", status) def _paginate(self, query): try: @@ -175,6 +218,14 @@ class ForumHandlers: return False return (datetime.now() - dt).total_seconds() < RECENT_SECONDS + def _get_subscribed_topics(self): + raw = self.fdb.get_setting("topic_subscriptions", "") + return [t.strip().lower() for t in raw.split(",") if t.strip()] + + def _get_subscribed_topics_str(self): + raw = self.fdb.get_setting("topic_subscriptions", "") + return raw + def _blocked_instances(self): raw = self.fdb.get_setting("blocked_instances", "") return set(h.strip() for h in raw.split(",") if h.strip()) @@ -204,6 +255,19 @@ class ForumHandlers: # --- Routes --- + def _status_bar(self): + topics = self._get_subscribed_topics() + topics_str = ", ".join(topics) if topics else "everything" + peer_count = len(self.fdb.get_synced_instances()) + filter_count = self.fdb.get_peer_filter_count() + return ( + f'
' + f'subscribed: {esc(topics_str)}' + f'{peer_count} peers' + f'{filter_count} filters' + f'
' + ) + def handle_list(self, query): page = self._paginate(query) tag = unquote(query.get("tag", [""])[0]).strip() @@ -224,12 +288,12 @@ class ForumHandlers: continue if self._is_new(r["created_at"]): new_count += 1 - badge = "[share]" if r["url"] else "[request]" - mute_badge = " [muted]" if is_muted else "" + badge = f'share' if r["url"] else f'request' + mute_label = " [muted]" if is_muted else "" tags_html = "" if r["tags"]: tag_links = " ".join( - f'[{esc(t.strip())}]' + f'{esc(t.strip())}' for t in r["tags"].split(",") if t.strip() ) tags_html = f' {tag_links}' @@ -237,7 +301,7 @@ class ForumHandlers: items += ( f'
  • ' f'
    ' - f'{badge}{mute_badge} ' + f'{badge}{mute_label} ' f'{esc(r["title"])}' f'{tags_html}' f'
    ' @@ -250,18 +314,19 @@ class ForumHandlers: f'
  • ' ) if not items: - items = "

    No threads yet.

    " + items = "

    no threads yet.

    " new_label = f" ({new_count} new)" if new_count else "" search_form = ( f'
    ' f'' f'
    ' ) - tag_label = f' — tag: {esc(tag)}' if tag else "" + tag_label = f' — {esc(tag)}' if tag else "" muted_link = f'show muted' if not show_muted else f'show all' page_url = f'/forum?q={esc(search)}&tag={esc(tag)}&muted=1' if show_muted else (f'/forum?q={esc(search)}&tag={esc(tag)}' if search or tag else '/forum') return self._respond( - f"

    forum{tag_label}

    " + f"

    forum{tag_label}

    " + f'{self._status_bar()}' f'
    ' f'{search_form}' f'
    ' @@ -270,26 +335,30 @@ class ForumHandlers: f'sync now' f'{muted_link}' f'
    ' - f"

    {total} threads{new_label}

    " + f"

    {total} threads{new_label}

    " f'' f"{self._page_nav(page, total, page_url)}" ) def handle_new_form(self, msg=""): return self._respond( - f"

    new thread

    " + f"

    new thread

    " f'
    ' f'{self._csrf_field()}' + f'' f'' f"max {MAX_TITLE_LENGTH} characters" + f'' f'' + f'' f'' f"max {MAX_BODY_LENGTH} characters" - f'' + f'' + f'' f'' f"
    " f"

    {msg}

    " - f'back' + f'
    back
    ' ) def handle_new_submit(self, body): @@ -323,28 +392,25 @@ class ForumHandlers: instance_hash = self.identity.hash.hex() if self.identity else "local" has_upvoted = self.fdb.has_upvoted(thread_id, instance_hash) - badge = "[share]" if thread["url"] else "[request]" + badge = f'share' if thread["url"] else f'request' url_html = "" if thread["url"]: url_html = ( f'

    {esc(thread["url"])}' - f' (+ save to my index)

    ' + f' (+ save)

    ' ) tags_html = "" if thread["tags"]: tag_links = " ".join( - f'[{esc(t.strip())}]' + f'{esc(t.strip())}' for t in thread["tags"].split(",") if t.strip() ) - tags_html = f'

    {tag_links}

    ' + tags_html = f'

    {tag_links}

    ' body_html = f"

    {esc(thread['body'])}

    " if thread["body"] else "" - mute_btn = ( - f'unmute' - if is_muted else - f'mute' - ) + mute_label = "unmute" if is_muted else "mute" + mute_href = f'/forum/unmute/{thread["id"]}' if is_muted else f'/forum/mute/{thread["id"]}' posts_html = "" for p in posts: @@ -352,18 +418,18 @@ class ForumHandlers: for word in p["body"].split(): w = word.strip().strip(",.!?;:") if w.startswith(("http://", "https://")): - save_links += ( - f' + save' - ) + save_links += f' + save' parent_ref = "" if p["parent_id"]: - parent_ref = f' ↪ reply' + parent_ref = f' ↪ reply' posts_html += ( - f'
    ' - f'{esc(self._author_str(p["author_name"], p["author_instance"]))}' + f'
    ' + f'' f'

    {esc(p["body"])}

    ' f'{save_links}' f'
    ' @@ -378,15 +444,16 @@ class ForumHandlers: f"" ) + upvote_label = "-1" if has_upvoted else "+1" return self._respond( - f"

    {badge} {esc(thread['title'])}

    " - f'

    ' + f"

    {badge} {esc(thread['title'])}

    " + f'

    ' f'by {esc(self._author_str(thread["author_name"], thread["author_instance"]))}' f'{self._block_link(thread["author_instance"])}' f' · {self._time_ago(thread["created_at"])}' f' · {thread["score"]} upvotes' - f' · {mute_btn}' - f' · {"-1" if has_upvoted else "+1"}' + f' · {mute_label}' + f' · {upvote_label}' f'{self._author_links(thread["id"], thread["author_instance"], instance_hash)}' f'

    ' f'{url_html}' @@ -428,7 +495,7 @@ class ForumHandlers: if thread["author_instance"] != instance_hash: return self._error(403) return self._respond( - f"

    edit thread

    " + f"

    edit thread

    " f'
    ' f'{self._csrf_field()}' f'' @@ -511,7 +578,7 @@ class ForumHandlers: def _peer_reports_html(self): counts = self.fdb.get_peer_block_counts() if not counts: - return "

    No peer reports yet.

    " + return "

    no peer reports yet

    " auto_blocked = set(h.strip() for h in self.fdb.get_setting("auto_blocked_instances", "").split(",") if h.strip()) blocked = self._blocked_instances() items = "" @@ -524,13 +591,24 @@ class ForumHandlers: blocked = self._blocked_instances() auto_blocked = set(h.strip() for h in self.fdb.get_setting("auto_blocked_instances", "").split(",") if h.strip()) peer_counts = self.fdb.get_peer_block_counts() + filters = self._keyword_filters() + filters_str = ", ".join(filters) if filters else "" + synced = self.fdb.get_synced_instances() + auto_discover = self.fdb.get_setting("forum_auto_discover", "1") + auto_discover_checked = " checked" if auto_discover == "1" else "" + auto_sync = self.fdb.get_setting("forum_auto_sync", "0") + auto_sync_checked = " checked" if auto_sync == "1" else "" + retention_days = self.fdb.get_setting("forum_retention_days", "30") + blocked_items = "" if blocked: for h in sorted(blocked): - label = "[auto] " if h in auto_blocked else "" - reports = f" ({peer_counts.get(h, 0)} peers)" if h in peer_counts else "" + label = "auto" if h in auto_blocked else "" + reports = f" ({peer_counts.get(h, 0)} reports)" if h in peer_counts else "" blocked_items += ( - f'
  • {label}{esc(h[:16])}...{reports} ' + f'
  • ' + f'{esc(h[:16])}...' + f'{" [" + label + "]" if label else ""}{reports} ' f'' f'{self._csrf_field()}' f'' @@ -539,12 +617,8 @@ class ForumHandlers: ) blocked_items = f"
      {blocked_items}
    " else: - blocked_items = "

    No instances blocked.

    " + blocked_items = "

    no instances blocked

    " - filters = self._keyword_filters() - filters_str = ", ".join(filters) if filters else "" - - synced = self.fdb.get_synced_instances() synced_items = "" for s in synced: synced_items += ( @@ -555,67 +629,78 @@ class ForumHandlers: f'
  • ' f'' ) - synced_items = f"" if synced_items else "

    No instances synced yet.

    " - - auto_discover = self.fdb.get_setting("forum_auto_discover", "1") - auto_discover_checked = " checked" if auto_discover == "1" else "" - auto_sync = self.fdb.get_setting("forum_auto_sync", "0") - auto_sync_checked = " checked" if auto_sync == "1" else "" - retention_days = self.fdb.get_setting("forum_retention_days", "30") + synced_items = f"" if synced_items else "

    no instances synced

    " return self._respond( - f"

    forum moderation

    " + f"

    moderation

    " f"

    {msg}

    " - f'

    sync now

    ' - f"

    auto-discovery

    " + f'{self._status_bar()}' + + f'
    ' + f'
    subscriptions
    ' + f'
    ' + f'{self._csrf_field()}' + f'' + f"only sync content matching these topics" + f'' + f"
    " + f'
    ' + + f'
    ' + f'
    settings
    ' + f'
    network behavior
    ' f'
    ' f'{self._csrf_field()}' f'" + f" auto-discover peers via announces" f'' f"
    " - f"

    auto-sync

    " f'
    ' f'{self._csrf_field()}' f'" + f" auto-sync every 5 minutes" f'' f"
    " - f"

    storage

    " f'
    ' f'{self._csrf_field()}' - f'' - f"Older threads are pruned automatically (default: 30). Set to 0 to keep everything." + f'' f'' f"
    " - f"

    blocked instances

    " - f"{blocked_items}" - f'
    ' - f'{self._csrf_field()}' - f'' - f'' - f"
    " - f"

    peer reports

    " - f"{self._peer_reports_html()}" - f"

    keyword filters

    " - f'
    ' - f'{self._csrf_field()}' - f'' - f'' - f"
    " - f"

    synced instances

    " - f"{synced_items}" - f"

    Instances are discovered automatically via mesh announces. " - f"You can also manually add a friend's instance hash to bootstrap.

    " + f'
    ' + + f'
    ' + f'
    network
    ' + f'
    {len(synced)} known peers
    ' + f'{synced_items}' f'
    ' f'{self._csrf_field()}' f'' f'' f'' f"
    " - f'
    ' - f'back to forum' + f'
    ' + + f'
    ' + f'
    moderation
    ' + f'{blocked_items}' + f'
    ' + f'{self._csrf_field()}' + f'' + f'' + f"
    " + f'
    peer reports
    ' + f'{self._peer_reports_html()}' + f'
    keyword filters
    ' + f'
    ' + f'{self._csrf_field()}' + f'' + f'' + f"
    " + f'
    ' + + f'' ) def handle_block(self, body): @@ -702,31 +787,122 @@ class ForumHandlers: self.fdb.remove_synced_instance(instance) return self.handle_moderation("Removed.") + def handle_topics(self, body): + topics = body.get("topics", [""])[0].strip() + self.fdb.set_setting("topic_subscriptions", topics) + return self.handle_moderation("Topic subscriptions saved.") + # --- Sync endpoint (called over RNS) --- def handle_sync_request(self, data): - """Handle incoming sync request from another forum instance.""" since = data.get("query", {}).get("since", [""])[0] if isinstance(data.get("query"), dict) else "" incoming_threads = data.get("threads", []) incoming_posts = data.get("posts", []) incoming_upvotes = data.get("upvotes", []) + peer_topics = data.get("my_topics", []) + peer_tag_cloud = data.get("my_tag_cloud", []) + peer_bloom_data = data.get("content_bloom") + peer_tag_bloom_data = data.get("tag_bloom") + peer_filter_table = data.get("filter_table", {}) + scoped_query_tag = data.get("scoped_query", "") blocked = self._blocked_instances() + my_topics = self._get_subscribed_topics() + from_hash = data.get("from_hash", "") + + # Store peer's tag bloom filter (Architecture B discovery) + if peer_tag_bloom_data and from_hash: + self.fdb.store_peer_filter( + peer_hash=from_hash, + bloom_bytes=bytes(peer_tag_bloom_data), + bloom_size=data.get("tag_bloom_size", 2048), + bloom_hashes=data.get("tag_bloom_hashes", 3), + tag_count=len(peer_topics), + ) + + # Merge filter table gossip (transitive peer discovery) + if peer_filter_table and from_hash: + for ph, entry in peer_filter_table.items(): + if isinstance(entry, dict): + bb = entry.get("bloom_bytes") + if bb: + self.fdb.store_peer_filter( + peer_hash=ph, + bloom_bytes=bytes(bb) if isinstance(bb, list) else bb, + bloom_size=entry.get("bloom_size", 2048), + bloom_hashes=entry.get("bloom_hashes", 3), + tag_count=entry.get("tag_count", 0), + ) + + # Handle scoped query: find peers whose bloom might contain the queried tag + scoped_query_results = [] + if scoped_query_tag and from_hash: + scoped_query_results = self.fdb.get_filtered_peers_by_tag(scoped_query_tag) + + # Store peer's topics for future routing (backward compat) + if peer_topics and from_hash: + self.fdb.set_setting(f"peer_topics_{from_hash}", ",".join(peer_topics)) + if peer_tag_cloud and from_hash: + self.fdb.set_setting(f"peer_tag_cloud_{from_hash}", json.dumps(peer_tag_cloud)) + + # Build our tag bloom filter to send back + peer_tag_bs = data.get("tag_bloom_size", 2048) + peer_tag_bh = data.get("tag_bloom_hashes", 3) + my_tag_bloom = BloomFilter.from_tags(my_topics or [], peer_tag_bs, peer_tag_bh) + + # Build our content bloom filter for dedup + our_existing = set() + for t in self.fdb.get_threads_by_topics(peer_topics) if peer_topics else []: + our_existing.add(t["id"]) + our_bloom = BloomFilter.from_items(list(our_existing), data.get("bloom_size", 2048) if data else 2048, data.get("bloom_hashes", 3) if data else 3) + + # Build filter table gossip from our stored filters + all_filters = self.fdb.get_all_filters() + filter_table_gossip = {} + for f in all_filters[:20]: + if f["peer_hash"] != from_hash: + filter_table_gossip[f["peer_hash"]] = { + "bloom_bytes": list(f["bloom_bytes"]), + "bloom_size": f["bloom_size"], + "bloom_hashes": f["bloom_hashes"], + "tag_count": f["tag_count"], + } + + # Parse peer's content bloom for dedup + peer_bloom = None + if peer_bloom_data: + peer_bloom = BloomFilter.from_bytes( + bytes(peer_bloom_data), + data.get("bloom_size", 2048), + data.get("bloom_hashes", 3), + ) + + # Merge incoming content, filtered by bloom and topics if incoming_threads: for t in incoming_threads: - if t.get("author_instance", "") not in blocked: + if t.get("author_instance", "") in blocked: + continue + if peer_bloom and peer_bloom.might_contain(t["id"]): + continue + if not my_topics: self.fdb.merge_thread(t) + else: + t_tags = [tag.strip().lower() for tag in t.get("tags", "").split(",") if tag.strip()] + if set(my_topics) & set(t_tags): + self.fdb.merge_thread(t) if incoming_posts: for p in incoming_posts: - if p.get("author_instance", "") not in blocked: - self.fdb.merge_post(p) + if p.get("author_instance", "") in blocked: + continue + if peer_bloom and peer_bloom.might_contain(p["id"]): + continue + self.fdb.merge_post(p) if incoming_upvotes: for uv in incoming_upvotes: self.fdb.merge_upvote(uv["thread_id"], uv["instance_hash"]) - # Record incoming peer blocks incoming_blocks = data.get("blocks", {}) - peer_hash = data.get("peer_hash", "") or data.get("from_hash", "") + peer_hash = data.get("peer_hash", "") or from_hash if incoming_blocks and peer_hash: for h in incoming_blocks.get("mine", []): if h and h not in blocked: @@ -735,13 +911,10 @@ class ForumHandlers: if h and h not in blocked: self.fdb.record_peer_block(peer_hash, h) - # Merge incoming retractions for r in data.get("retractions", []): if r.get("id") and r.get("type") and r.get("author") and r.get("at"): self.fdb.merge_retraction(r["id"], r["type"], r["author"], r["at"]) - # Auto-discover the peer that synced with us and their known peers - from_hash = data.get("from_hash", "") if from_hash and from_hash not in blocked: self.fdb.add_known_peer(from_hash) for peer_hash in data.get("known_peers", []): @@ -750,18 +923,34 @@ class ForumHandlers: my_blocks = list(blocked) my_peer_blocks = self.fdb.get_peer_block_list() + my_tag_cloud = self.fdb.get_tag_cloud(50) + my_tag_list = [t for t, _ in my_tag_cloud] + threads, posts, upvote_threads = [], [], [] if since: - ts, posts_list, up_list = self.fdb.get_new_content(since) - threads = [dict(r) for r in ts] - posts = [dict(r) for r in posts_list] - upvote_threads = up_list + if peer_topics: + rows = self.fdb.get_threads_by_topics(peer_topics, since=since) + threads = [dict(r) for r in rows] + tids = [r["id"] for r in rows] + posts = [dict(p) for p in self.fdb.get_posts_by_thread_ids(tids)] + uv_rows = self.fdb.get_new_upvotes_since(since, tids) + upvote_threads = uv_rows + else: + ts, posts_list, up_list = self.fdb.get_new_content(since) + threads = [dict(r) for r in ts] + posts = [dict(r) for r in posts_list] + upvote_threads = up_list + + known_peers = [h for h in self.fdb.get_all_known_hashes() if h != from_hash] + peer_topics_map = {} + for ph in known_peers[:100]: + pt = self.fdb.get_setting(f"peer_topics_{ph}", "") + if pt: + peer_topics_map[ph] = [t.strip() for t in pt.split(",") if t.strip()] retracted = [{"id": cid, "type": ct, "author": ai, "at": ra} for cid, ct, ai, ra in self.fdb.get_raw_retractions()] - known_peers = [h for h in self.fdb.get_all_known_hashes() if h != from_hash] - return { "status": 200, "content_type": "application/json", @@ -772,6 +961,18 @@ class ForumHandlers: "blocks": {"mine": my_blocks, "peers": my_peer_blocks}, "retractions": retracted, "known_peers": known_peers, + "peer_topics": my_tag_list, + "peer_tag_cloud": my_tag_cloud, + "content_bloom": list(our_bloom.bytes), + "bloom_size": data.get("bloom_size", 2048), + "bloom_hashes": data.get("bloom_hashes", 3), + "peer_topics_map": peer_topics_map, + # Architecture B additions + "tag_bloom": list(my_tag_bloom.bytes), + "tag_bloom_size": peer_tag_bs, + "tag_bloom_hashes": peer_tag_bh, + "filter_table": filter_table_gossip, + "scoped_query_results": scoped_query_results, }), "headers": {}, } @@ -832,7 +1033,7 @@ class ForumHandlers: elif method == "POST": if not self._check_csrf(body): return self._with_csrf( - self._respond("

    403 Forbidden

    ", status=403), csrf_token + self._respond("

    403 Forbidden

    ", status=403), csrf_token ) if sub == "/new": return self._with_csrf(self.handle_new_submit(body), csrf_token) @@ -867,6 +1068,8 @@ class ForumHandlers: return self._with_csrf(self.handle_storage(body), csrf_token) elif sub == "/auto_sync": return self._with_csrf(self.handle_auto_sync(body), csrf_token) + elif sub == "/topics": + return self._with_csrf(self.handle_topics(body), csrf_token) return self._with_csrf(self._error(404), csrf_token) diff --git a/tinyweb_forum/sync.py b/tinyweb_forum/sync.py index 99e1985..a9d544d 100644 --- a/tinyweb_forum/sync.py +++ b/tinyweb_forum/sync.py @@ -4,15 +4,22 @@ import threading import time import RNS +from tinyweb_forum.bloom import BloomFilter + FORUM_APP = "tinyweb-forum" -SYNC_INTERVAL = 300 # 5 minutes +SYNC_INTERVAL = 300 REQUEST_TIMEOUT = 60 -GOSSIP_FANOUT = 20 # random peers to sync per cycle +GOSSIP_FANOUT = 20 +BLOOM_SIZE = 2048 +BLOOM_HASHES = 3 +TAG_BLOOM_SIZE = 2048 +TAG_BLOOM_HASHES = 3 +FILTER_TABLE_GOSSIP = 20 +MAX_PEER_FAILURES = 5 +FILTER_TABLE_TTL_DAYS = 7 class _ForumAnnounceHandler: - """Receives announces from other forum instances and auto-discovers them.""" - aspect_filter = FORUM_APP receive_path_responses = False @@ -39,6 +46,26 @@ class ForumSync: self._running = False self._thread = None + def _get_subscribed_topics(self): + raw = self.fdb.get_setting("topic_subscriptions", "") + return [t.strip().lower() for t in raw.split(",") if t.strip()] + + def _build_tag_bloom(self, topics=None): + if topics is None: + topics = self._get_subscribed_topics() + return BloomFilter.from_tags(topics or [], TAG_BLOOM_SIZE, TAG_BLOOM_HASHES) + + def _topics_overlap(self, my_topics, their_topics): + if not my_topics or not their_topics: + return True + return bool(set(my_topics) & set(their_topics)) + + def _topics_overlap_bloom(self, my_topics, peer_bloom): + """Check overlap using bloom filter instead of topic list.""" + if not my_topics or peer_bloom is None: + return True + return any(peer_bloom.might_contain(t) for t in my_topics) + def start(self): self.destination = RNS.Destination( self.identity, @@ -66,23 +93,48 @@ class ForumSync: self.fdb.set_setting("forum_auto_sync", "1" if enabled else "0") if enabled: self._start_sync_loop() - else: - pass # current cycle finishes, no new one starts def sync_now(self): - """Run one sync cycle immediately. Returns count of peers synced.""" instances = self.fdb.get_synced_instances() random.shuffle(instances) count = 0 + my_topics = self._get_subscribed_topics() + my_tag_bloom = self._build_tag_bloom(my_topics) + for inst in instances[:GOSSIP_FANOUT]: + if not self._running: + break + pfilter = self.fdb.get_peer_filter(inst["instance_hash"]) + if pfilter: + pbloom = BloomFilter.from_bytes(pfilter["bloom_bytes"], pfilter["bloom_size"], pfilter["bloom_hashes"]) + if not self._topics_overlap_bloom(my_topics, pbloom): + continue + else: + peer_topics = self._peer_tag_topics(inst["instance_hash"]) + if not self._topics_overlap(my_topics, peer_topics): + continue + try: + self._sync_with(inst["instance_hash"], my_topics, my_tag_bloom) + count += 1 + except Exception as e: + self.fdb.record_sync_result(inst["instance_hash"], False) + print(f"[forum] sync error with {inst['instance_hash'][:16]}: {e}") + return count + + def scoped_query(self, tag): + """Ask known peers if they know anyone with the given tag.""" + tag = tag.strip().lower() + instances = self.fdb.get_synced_instances() + random.shuffle(instances) + results = set() for inst in instances[:GOSSIP_FANOUT]: if not self._running: break try: - self._sync_with(inst["instance_hash"]) - count += 1 - except Exception as e: - print(f"[forum] sync error with {inst['instance_hash'][:16]}: {e}") - return count + matches = self._scoped_query_peer(inst["instance_hash"], tag) + results.update(matches) + except Exception: + continue + return list(results) def _start_sync_loop(self): if self._thread and self._thread.is_alive(): @@ -120,16 +172,25 @@ class ForumSync: while self._running: try: instances = self.fdb.get_synced_instances() - # Gossip: sync with random subset for scaling - # If <= GOSSIP_FANOUT peers, sync with all (current behavior) - # If more, sync with random FANOUT per cycle — content spreads epidemically + my_topics = self._get_subscribed_topics() + my_tag_bloom = self._build_tag_bloom(my_topics) random.shuffle(instances) for inst in instances[:GOSSIP_FANOUT]: if not self._running: break + pfilter = self.fdb.get_peer_filter(inst["instance_hash"]) + if pfilter: + pbloom = BloomFilter.from_bytes(pfilter["bloom_bytes"], pfilter["bloom_size"], pfilter["bloom_hashes"]) + if not self._topics_overlap_bloom(my_topics, pbloom): + continue + else: + peer_topics = self._peer_tag_topics(inst["instance_hash"]) + if not self._topics_overlap(my_topics, peer_topics): + continue try: - self._sync_with(inst["instance_hash"]) + self._sync_with(inst["instance_hash"], my_topics, my_tag_bloom) except Exception as e: + self.fdb.record_sync_result(inst["instance_hash"], False) print(f"[forum] sync error with {inst['instance_hash'][:16]}: {e}") except Exception: pass @@ -141,12 +202,99 @@ class ForumSync: self.fdb.prune_old_content(days) except Exception: pass + try: + self.fdb.prune_peer_filters(FILTER_TABLE_TTL_DAYS) + except Exception: + pass + try: + for dead in self.fdb.get_dead_peers(MAX_PEER_FAILURES): + self.fdb.remove_synced_instance(dead) + except Exception: + pass for _ in range(SYNC_INTERVAL): if not self._running: return time.sleep(1) - def _sync_with(self, instance_hash): + def _peer_tag_topics(self, instance_hash): + raw = self.fdb.get_setting(f"peer_topics_{instance_hash}", "") + return [t.strip().lower() for t in raw.split(",") if t.strip()] if raw else [] + + def _store_peer_filter(self, peer_hash, data): + tag_bloom_data = data.get("tag_bloom") + if tag_bloom_data: + self.fdb.store_peer_filter( + peer_hash=peer_hash, + bloom_bytes=bytes(tag_bloom_data), + bloom_size=data.get("tag_bloom_size", TAG_BLOOM_SIZE), + bloom_hashes=data.get("tag_bloom_hashes", TAG_BLOOM_HASHES), + tag_count=len(data.get("peer_topics", [])), + ) + + def _merge_filter_table_gossip(self, filter_table): + if not filter_table: + return + for ph, entry in filter_table.items(): + if isinstance(entry, dict): + bb = entry.get("bloom_bytes") + if bb: + self.fdb.store_peer_filter( + peer_hash=ph, + bloom_bytes=bytes(bb) if isinstance(bb, list) else bb, + bloom_size=entry.get("bloom_size", TAG_BLOOM_SIZE), + bloom_hashes=entry.get("bloom_hashes", TAG_BLOOM_HASHES), + tag_count=entry.get("tag_count", 0), + ) + + def _scoped_query_peer(self, peer_hash, tag): + dest_hash = bytes.fromhex(peer_hash) + if not RNS.Transport.has_path(dest_hash): + RNS.Transport.request_path(dest_hash) + elapsed = 0 + while not RNS.Transport.has_path(dest_hash) and elapsed < 15: + time.sleep(0.5) + elapsed += 0.5 + if not RNS.Transport.has_path(dest_hash): + return [] + server_identity = RNS.Identity.recall(dest_hash) + if server_identity is None: + return [] + destination = RNS.Destination( + server_identity, + RNS.Destination.OUT, + RNS.Destination.SINGLE, + FORUM_APP, + ) + link = RNS.Link(destination) + elapsed = 0 + while link.status == RNS.Link.PENDING and elapsed < 15: + time.sleep(0.25) + elapsed += 0.25 + if link.status != RNS.Link.ACTIVE: + return [] + try: + my_hash = self.identity.hash.hex() if self.identity else "local" + request_data = { + "scoped_query": tag, + "from_hash": my_hash, + } + receipt = link.request("/forum", data=request_data, timeout=REQUEST_TIMEOUT) + elapsed = 0 + done = (RNS.RequestReceipt.READY, RNS.RequestReceipt.DELIVERED, RNS.RequestReceipt.FAILED) + while receipt.get_status() not in done and elapsed < REQUEST_TIMEOUT: + time.sleep(0.1) + elapsed += 0.1 + if receipt.get_status() in (RNS.RequestReceipt.READY, RNS.RequestReceipt.DELIVERED): + resp = receipt.get_response() + if isinstance(resp, dict) and resp.get("status") == 200: + data = json.loads(resp["body"]) + return data.get("scoped_query_results", []) + return [] + finally: + link.teardown() + + def _sync_with(self, instance_hash, my_topics=None, my_tag_bloom=None): + success = False dest_hash = bytes.fromhex(instance_hash) if not RNS.Transport.has_path(dest_hash): RNS.Transport.request_path(dest_hash) @@ -155,10 +303,12 @@ class ForumSync: time.sleep(0.5) elapsed += 0.5 if not RNS.Transport.has_path(dest_hash): + self.fdb.record_sync_result(instance_hash, False) return server_identity = RNS.Identity.recall(dest_hash) if server_identity is None: + self.fdb.record_sync_result(instance_hash, False) return destination = RNS.Destination( @@ -175,6 +325,7 @@ class ForumSync: elapsed += 0.25 if link.status != RNS.Link.ACTIVE: + self.fdb.record_sync_result(instance_hash, False) return try: @@ -186,14 +337,50 @@ class ForumSync: break since = last_sync.replace(" ", "T") if last_sync else "" + my_hash = self.identity.hash.hex() if self.identity else "local" + + if my_topics is None: + my_topics = self._get_subscribed_topics() + if my_tag_bloom is None: + my_tag_bloom = self._build_tag_bloom(my_topics) + + # Build filter table gossip: share a subset of our known peer filters + all_filters = self.fdb.get_all_filters() + filter_table_gossip = {} + for f in all_filters[:FILTER_TABLE_GOSSIP]: + if f["peer_hash"] != instance_hash and f["peer_hash"] != my_hash: + filter_table_gossip[f["peer_hash"]] = { + "bloom_bytes": list(f["bloom_bytes"]), + "bloom_size": f["bloom_size"], + "bloom_hashes": f["bloom_hashes"], + "tag_count": f["tag_count"], + } threads, posts = [], [] upvotes = [] + my_tag_cloud = self.fdb.get_tag_cloud(50) + my_tag_list = [t for t, _ in my_tag_cloud] + + # Only build content if we have a since timestamp (incremental sync) if since: - ts, ps, uv = self.fdb.get_new_content(since) - threads = [dict(r) for r in ts] - posts = [dict(r) for r in ps] - upvotes = [{"thread_id": tid, "instance_hash": instance_hash} for tid in uv] + if my_topics: + rows = self.fdb.get_threads_by_topics(my_topics, since=since) + threads = [dict(r) for r in rows] + tids = [r["id"] for r in rows] + posts = [dict(p) for p in self.fdb.get_posts_by_thread_ids(tids)] + uv_rows = self.fdb.get_new_upvotes_since(since, tids) + upvotes = [{"thread_id": tid, "instance_hash": instance_hash} for tid in uv_rows] + else: + ts, ps, uv = self.fdb.get_new_content(since) + threads = [dict(r) for r in ts] + posts = [dict(r) for r in ps] + upvotes = [{"thread_id": tid, "instance_hash": instance_hash} for tid in uv] + + # Build content bloom for dedup + existing_ids = set(t["id"] for t in threads) + for t in self.fdb.get_threads_by_topics(my_topics) if my_topics else []: + existing_ids.add(t["id"]) + content_bloom = BloomFilter.from_items(list(existing_ids), BLOOM_SIZE, BLOOM_HASHES) my_blocks = [h.strip() for h in self.fdb.get_setting("blocked_instances", "").split(",") if h.strip()] my_peer_blocks = self.fdb.get_peer_block_list() @@ -201,8 +388,12 @@ class ForumSync: retracted = [{"id": cid, "type": ct, "author": ai, "at": ra} for cid, ct, ai, ra in self.fdb.get_raw_retractions()] - my_hash = self.identity.hash.hex() if self.identity else "local" known_peers = [h for h in self.fdb.get_all_known_hashes() if h != instance_hash and h != my_hash] + peer_topics_map = {} + for ph in known_peers[:100]: + pt = self.fdb.get_setting(f"peer_topics_{ph}", "") + if pt: + peer_topics_map[ph] = [t.strip() for t in pt.split(",") if t.strip()] request_data = { "query": {"since": [since]} if since else {}, @@ -213,6 +404,17 @@ class ForumSync: "blocks": {"mine": my_blocks, "peers": my_peer_blocks}, "retractions": retracted, "known_peers": known_peers, + "my_topics": my_topics, + "my_tag_cloud": my_tag_cloud, + "content_bloom": list(content_bloom.bytes), + "bloom_size": BLOOM_SIZE, + "bloom_hashes": BLOOM_HASHES, + "peer_topics_map": peer_topics_map, + # Architecture B additions + "tag_bloom": list(my_tag_bloom.bytes), + "tag_bloom_size": TAG_BLOOM_SIZE, + "tag_bloom_hashes": TAG_BLOOM_HASHES, + "filter_table": filter_table_gossip, } receipt = link.request("/forum", data=request_data, timeout=REQUEST_TIMEOUT) @@ -229,16 +431,82 @@ class ForumSync: data = json.loads(resp["body"]) except (json.JSONDecodeError, KeyError): data = {} + success = True + + # Store peer's tag bloom filter + self._store_peer_filter(instance_hash, data) + + # Merge filter table gossip from peer + self._merge_filter_table_gossip(data.get("filter_table", {})) + + # Store peer topics (backward compat) + peer_topics = data.get("peer_topics", []) + if peer_topics: + self.fdb.set_setting(f"peer_topics_{instance_hash}", ",".join(peer_topics)) + peer_tag_cloud = data.get("peer_tag_cloud", []) + if peer_tag_cloud: + self.fdb.set_setting(f"peer_tag_cloud_{instance_hash}", json.dumps(peer_tag_cloud)) + + # Check tag overlap using bloom filter my_blocks = set(h.strip() for h in self.fdb.get_setting("blocked_instances", "").split(",") if h.strip()) + + peer_tag_bloom_data = data.get("tag_bloom") + peer_tag_bloom = None + if peer_tag_bloom_data: + peer_tag_bloom = BloomFilter.from_bytes( + bytes(peer_tag_bloom_data), + data.get("tag_bloom_size", TAG_BLOOM_SIZE), + data.get("tag_bloom_hashes", TAG_BLOOM_HASHES), + ) + + # If no tag overlap, skip content sync but still store the filter + if my_topics and peer_tag_bloom is not None: + if not self._topics_overlap_bloom(my_topics, peer_tag_bloom): + success = True + now = time.strftime("%Y-%m-%dT%H:%M:%S") + self.fdb.set_last_sync(instance_hash, now) + self.fdb.record_sync_result(instance_hash, True) + return + + # Fall back to topic list overlap check if no bloom + if peer_tag_bloom is None: + incoming_topics = data.get("peer_topics", []) + if my_topics and incoming_topics: + if not self._topics_overlap(my_topics, incoming_topics): + success = True + now = time.strftime("%Y-%m-%dT%H:%M:%S") + self.fdb.set_last_sync(instance_hash, now) + self.fdb.record_sync_result(instance_hash, True) + return + + # Check content bloom for dedup + peer_bloom_data = data.get("content_bloom") + peer_bloom = None + if peer_bloom_data: + peer_bloom = BloomFilter.from_bytes( + bytes(peer_bloom_data), + data.get("bloom_size", BLOOM_SIZE), + data.get("bloom_hashes", BLOOM_HASHES), + ) + for t in data.get("threads", []): if t.get("author_instance", "") not in my_blocks: - self.fdb.merge_thread(t) + if peer_bloom and peer_bloom.might_contain(t["id"]): + continue + if not my_topics: + self.fdb.merge_thread(t) + else: + t_tags = [tag.strip().lower() for tag in t.get("tags", "").split(",") if tag.strip()] + if set(my_topics) & set(t_tags): + self.fdb.merge_thread(t) for p in data.get("posts", []): if p.get("author_instance", "") not in my_blocks: + if peer_bloom and peer_bloom.might_contain(p["id"]): + continue self.fdb.merge_post(p) for tid in data.get("upvote_threads", []): self.fdb.merge_upvote(tid, instance_hash) - # Gossip blocks from peer + peer_blocks = data.get("blocks", {}) for h in peer_blocks.get("mine", []): if h and h not in my_blocks and instance_hash: @@ -247,18 +515,25 @@ class ForumSync: if h and h not in my_blocks and instance_hash: self.fdb.record_peer_block(instance_hash, h) self._apply_peer_blocks() - # Merge incoming retractions + for r in data.get("retractions", []): if r.get("id") and r.get("type") and r.get("author") and r.get("at"): self.fdb.merge_retraction(r["id"], r["type"], r["author"], r["at"]) - # Discover new peers from gossip + for peer_hash in data.get("known_peers", []): if peer_hash and peer_hash != my_hash and peer_hash != instance_hash: self.fdb.add_known_peer(peer_hash) + + my_topics_set = set(my_topics) + for ph, pt in data.get("peer_topics_map", {}).items(): + if ph and ph != my_hash and ph != instance_hash: + pt_set = set(t.lower() for t in pt) + if not my_topics or my_topics_set & pt_set: + self.fdb.add_known_peer(ph) + now = time.strftime("%Y-%m-%dT%H:%M:%S") self.fdb.set_last_sync(instance_hash, now) - else: - pass + self.fdb.record_sync_result(instance_hash, success) finally: link.teardown()