From 497079c7a74ee867e06ae457c5e64904b5d5f160 Mon Sep 17 00:00:00 2001
From: blankie
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''
)
- 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"{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"{msg}
"
- f'back '
+ f''
)
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(self._author_str(p["author_name"], p["author_instance"]))} '
f'{self._block_link(p["author_instance"])}'
f' · {self._time_ago(p["created_at"])}{parent_ref}'
- f'{" · " + self._post_retract_link(thread["id"], p["id"]) if p["author_instance"] == instance_hash else ""}'
+ f'{" · " + self._post_retract_link(thread["id"], p["id"]) if p["author_instance"] == instance_hash else ""}'
+ 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""
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'remove '
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'save '
+ f" "
+ f'
'
+
+ f''
+ f'
settings
'
+ f'
network behavior
'
f'
'
f'{self._csrf_field()}'
f' '
- f" automatically discover other forum instances on the mesh "
+ f" auto-discover peers via announces"
f'save '
f" "
- f"
auto-sync "
f'
'
f'{self._csrf_field()}'
f' '
- f" automatically sync content every 5 minutes "
+ f" auto-sync every 5 minutes"
f'save '
f" "
- f"
storage "
f'
'
f'{self._csrf_field()}'
- f'Keep threads for '
- f' days '
- f"Older threads are pruned automatically (default: 30). Set to 0 to keep everything. "
+ f'keep threads for '
+ f' '
+ f' days (0 = keep everything) '
f'save '
f" "
- f"
blocked instances "
- f"{blocked_items}"
- f'
'
- f'{self._csrf_field()}'
- f' '
- f'block '
- f" "
- f"
peer reports "
- f"{self._peer_reports_html()}"
- f"
keyword filters "
- f'
'
- f'{self._csrf_field()}'
- f' '
- f'save '
- 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'add '
f" "
- f'
'
- f'
back to forum '
+ f'
'
+
+ f''
+ f'
moderation
'
+ f'{blocked_items}'
+ f'
'
+ f'{self._csrf_field()}'
+ f' '
+ f'block '
+ f" "
+ f'
peer reports
'
+ f'{self._peer_reports_html()}'
+ f'
keyword filters
'
+ f'
'
+ f'{self._csrf_field()}'
+ f' '
+ f'save '
+ 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()