Repository navigation
impl learnings and add configuration options #23
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
eabc354
a6a6a79
b58fdfe
391abcc
a0c73fa
93b909e
1de3c6e
8afe88f
10303f4
bcb94c6
7b210c3
341ff8c
24c8bb5
66158a5
44163f8
a9d0ce1
4d7bad7
1ba92ff
68557ad
44f2697
2d6d793
07405ac
b410eca
fa81706
7701583
228de6e
630e8a1
6d82584
a928778
e29f18a
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -17,7 +17,9 @@ Contents: | |
| Components: | ||
| - iroh endpoint and gossip | ||
| - Auto-discovery (bootstrap loop) | ||
| - Publisher (background task) | ||
| - Publisher (background actor) | ||
| - Bubble merge (background actor) | ||
| - Message overlap merge (background actor) | ||
| - DHT client (mutable records) | ||
| - Crypto (signing, encryption, secret rotation) | ||
|
|
||
|
|
@@ -35,6 +37,8 @@ flowchart LR | |
| subgraph AutoDiscovery | ||
| B[Bootstrap Loop] | ||
| P[Publisher] | ||
| BM[Bubble Merge] | ||
| MO[Message Overlap Merge] | ||
| end | ||
|
|
||
| subgraph DHT | ||
|
|
@@ -44,16 +48,20 @@ flowchart LR | |
| A --> E --> G | ||
| G <---> B | ||
| G <---> P | ||
| G <---> BM | ||
| G <---> MO | ||
| B <--> D | ||
| P <--> D | ||
| BM <--> D | ||
| MO <--> D | ||
| ``` | ||
|
|
||
| Node lifecycle: | ||
| - Start iroh endpoint | ||
| - Start gossip | ||
| - Auto-discovery: | ||
| - Join topic, attempt bootstrap, connect | ||
| - Spawn publisher on success | ||
| - Spawn publisher, bubble merge, and message overlap merge actors on success | ||
|
|
||
| State machine: | ||
|
|
||
|
|
@@ -64,7 +72,7 @@ stateDiagram-v2 | |
| Discovering --> Joining | ||
| Joining --> Joined | ||
| Joined --> Publishing | ||
| Publishing --> Joined : backoff or success loop | ||
| Publishing --> Joined : interval + jitter | ||
| Discovering --> Discovering : retry/jitter | ||
| Joining --> Discovering : no peers | ||
| ``` | ||
|
|
@@ -82,25 +90,28 @@ sequenceDiagram | |
| participant Gossip | ||
|
|
||
| Node->>Gossip: subscribe(topic_hash) | ||
| Node->>DHT: get_mutable(signing_pub, salt, 10s) | ||
| Note over Node: optionally publish startup record | ||
| Node->>DHT: get_mutable(signing_pub, salt, 10s timeout) | ||
| DHT-->>Node: encrypted records (0..N) | ||
| Node->>Node: decrypt, verify, filter(not self) | ||
| alt candidates exist | ||
| loop each candidate | ||
| Node->>Gossip: join_peers([node_id]) | ||
| Node->>Node: sleep 100ms | ||
| Node->>Gossip: join_peers([pub_key]) | ||
| Node->>Node: sleep per_peer_join_settle_time (100ms) | ||
| Gossip-->>Node: NeighborUp? | ||
| end | ||
| Node->>Node: final wait 500ms | ||
| Node->>Node: final wait join_confirmation_wait_time (500ms) | ||
| else no candidates | ||
| Node->>Node: maybe publish own (rate-limited) | ||
| Node->>Node: sleep no_peers_retry_interval (1500ms) | ||
| end | ||
| Node->>Node: joined? if yes, spawn publisher | ||
| Node->>Node: joined? if yes, spawn publisher + merge actors | ||
| ``` | ||
|
|
||
| Key points: | ||
| - First iteration: also check previous unix minute. | ||
| - Pacing avoids bursts and “bubbles.” | ||
| - First iteration: optionally check older records first (`check_older_records_first_on_startup`). | ||
| - Both `unix_minute` and `unix_minute - 1` records are always fetched. | ||
| - Pacing avoids bursts and "bubbles." | ||
| - Keep trying until joined. | ||
|
|
||
| Pseudocode: | ||
|
|
@@ -109,23 +120,23 @@ Pseudocode: | |
| loop: | ||
| if joined(): return sender, receiver | ||
|
|
||
| minute = first_attempt ? -1 : 0 | ||
| recs = get_unix_minute_records(minute) | ||
| minute = first_attempt && check_last_minute_first ? -1 : 0 | ||
| recs = get_records(unix_minute(minute) - 1) + get_records(unix_minute(minute)) | ||
|
|
||
| if recs.is_empty(): | ||
| maybe_publish_this_minute() | ||
| sleep(100ms) | ||
| sleep(no_peers_retry_interval = 1500ms) | ||
| continue | ||
|
|
||
| for peer in extract_bootstrap_nodes(recs): | ||
| if joined(): break | ||
| join_peer(peer) | ||
| sleep(100ms) | ||
| sleep(per_peer_join_settle_time = 100ms) | ||
|
|
||
| sleep(500ms) | ||
| sleep(join_confirmation_wait_time = 500ms) | ||
| if joined(): return | ||
| maybe_publish_this_minute() | ||
| sleep(100ms) | ||
| sleep(discovery_poll_interval = 2000ms) | ||
| ``` | ||
|
|
||
| ## Publishing | ||
|
|
@@ -136,36 +147,36 @@ Flow: | |
|
|
||
| ```mermaid | ||
| flowchart TD | ||
| A[Start Cycle] --> B[Get minute=now] | ||
| B --> C[Discover existing records] | ||
| C --> D[Filter active participants] | ||
| D --> E{>= 10 active?} | ||
| E -- Yes --> F[Stop rate-limited] | ||
| A[Tick] --> B[Get minute=now] | ||
| B --> C[Get existing records] | ||
| C --> E{>= max_bootstrap_records = 5?} | ||
| E -- Yes --> F[Skip - rate-limited] | ||
| E -- No --> G[Build record: peers + msg hashes] | ||
| G --> H[Sign + Encrypt] | ||
| H --> I[Publish with retries + jitter] | ||
| I --> J[Return records, including own on success] | ||
| H --> I[Publish to DHT] | ||
| I --> J[Reset ticker: base_interval + random jitter] | ||
| ``` | ||
|
|
||
| Pseudocode: | ||
|
|
||
| ```text | ||
| records = get_unix_minute_records(now) | ||
| active = filter_active(records) | ||
| if active.len >= 10: return records | ||
|
|
||
| rec = make_record(neighbors(<=5), last_hashes(<=5)) | ||
| enc = encrypt(sign(rec)) | ||
| publish_with_retry(enc, retries=3, jitter=0..2000ms) | ||
| return records + [rec_if_success] | ||
| // Publisher actor loop (interval: base_interval + random jitter) | ||
| on tick: | ||
| records = get_records(unix_minute(0)) | ||
| if records.len >= max_bootstrap_records(5): return | ||
|
|
||
| rec = make_record(neighbors(<=5), last_hashes(<=5)) | ||
| enc = encrypt(sign(rec)) | ||
| publish(enc) | ||
| reset_ticker(base_interval + random(0, max_jitter)) | ||
| ``` | ||
|
|
||
| ## Bubble detection and merging | ||
|
|
||
| Signal 1: small cluster \(neighbors < 4\). | ||
| Signal 1: small cluster \(neighbors < min\_neighbors, default 4\). | ||
| - Extract peer ids from discovered records. | ||
| - Exclude zeros, self, current neighbors. | ||
| - Join up to MAX_JOIN_PEERS_COUNT. | ||
| - Join up to max_join_peer_count (default 4). | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Documented parameter name doesn't match the code. The docs call the bubble-merge join cap 📝 Proposed fix-- Join up to max_join_peer_count (default 4).
+- Join up to `max_join_peers` (default 4).
@@
-- `max_join_peer_count` (default 4)
+- `max_join_peers` (default 4)Also applies to: 273-280 🤖 Prompt for AI Agents |
||
|
|
||
| Signal 2: non-overlapping message sets. | ||
| - Compare local last_message_hashes with others. | ||
|
|
@@ -176,12 +187,12 @@ Decision graph: | |
|
|
||
| ```mermaid | ||
| flowchart LR | ||
| A[Post-Publish Records] --> B{neighbors < 4?} | ||
| A[Tick] --> B{neighbors < min_neighbors?} | ||
| B -- Yes --> C[Join peers from records] | ||
| B -- No --> D{local_msgs >= 1?} | ||
| D -- No --> E[Sleep random 0..60s] | ||
| D -- No --> E[Sleep until next tick] | ||
| D -- Yes --> F{overlap with others?} | ||
| F -- No --> G[Join from non-overlap records] | ||
| F -- No --> G[Join from non-overlapping records] | ||
| F -- Yes --> E | ||
| ``` | ||
|
|
||
|
|
@@ -190,9 +201,8 @@ flowchart LR | |
| Record (summary): | ||
| - topic hash (32) | ||
| - unix_minute (u64) | ||
| - node_id (publisher) | ||
| - active_peers[5] (node ids) | ||
| - last_message_hashes[5] | ||
| - pub_key (publisher ed25519 public key) | ||
| - content (serialized GossipRecordContent: active_peers + last_message_hashes) | ||
| - signature (64) | ||
|
|
||
| EncryptedRecord: | ||
|
|
@@ -206,16 +216,22 @@ classDiagram | |
| class Record { | ||
| +topic: [u8;32] | ||
| +unix_minute: u64 | ||
| +node_id: [u8;32] | ||
| +pub_key: [u8;32] | ||
| +content: GossipRecordContent | ||
| +signature: [u8;64] | ||
| } | ||
|
|
||
| class GossipRecordContent { | ||
| +active_peers: [[u8;32];5] | ||
| +last_message_hashes: [[u8;32];5] | ||
| +signature: [u8;64] | ||
| } | ||
|
|
||
| class EncryptedRecord { | ||
| +encrypted_record: Vec<u8> | ||
| +encrypted_decryption_key: Vec<u8> | ||
| } | ||
|
|
||
| Record --* GossipRecordContent : content deserializes to | ||
| ``` | ||
|
|
||
| Key derivation: | ||
|
|
@@ -225,7 +241,9 @@ flowchart TD | |
| T[topic_hash] --> A[SHA512 topic+minute] | ||
| M[unix_minute] --> A | ||
| A --> S[signing_keypair seed -> Ed25519] | ||
| A --> L[salt = first 32 bytes] | ||
|
|
||
| T --> L["salt = SHA512('salt' + topic + minute)[..32]"] | ||
| M --> L | ||
|
|
||
| T --> R[secret_rotation topic,minute,initial_secret_hash] | ||
| M --> R | ||
|
|
@@ -239,23 +257,24 @@ flowchart TD | |
| - Decrypt/verify failure: | ||
| - Drop record; proceed. | ||
| - Publish failure: | ||
| - Exponential backoff (1..60 s), then retry. | ||
| - DHT layer retries with jittered intervals (3 retries, 5s base + 0-10s jitter). | ||
| - Join failure: | ||
| - Continue to next peer; final 500 ms wait; loop. | ||
| - Continue to next peer; final 500ms wait; loop. | ||
|
|
||
| ## Tuning | ||
|
|
||
| - Per-minute cap \(N_{active} \ge 10\) gates publishing. | ||
| - Pacing (100 ms) reduces bursts. | ||
| - Backoff (1..60 s) stabilizes DHT load. | ||
| - Per-minute cap \(records \ge max\_bootstrap\_records, default 5\) gates publishing. | ||
| - Per-peer pacing (100ms) reduces bursts. | ||
| - No-peers retry (1500ms) and discovery poll (2000ms) stabilize DHT load. | ||
| - Message window size (5 peers, 5 hashes) is a trade-off: | ||
| - Larger window = better visibility, larger records. | ||
| - Smaller window = lower bandwidth, less overlap detection. | ||
|
|
||
| Parameters: | ||
| - MAX_BOOTSTRAP_RECORDS | ||
| - MAX_JOIN_PEERS_COUNT | ||
| - DHT timeout | ||
| - Retry count and jitter | ||
| - Join pacing and final wait | ||
| - Publisher backoff and success jitter | ||
| Parameters (all configurable): | ||
| - `max_bootstrap_records` (default 5) | ||
| - `max_join_peer_count` (default 4) | ||
| - `min_neighbors` for bubble merge (default 4) | ||
| - DHT timeouts, retry count, and jitter | ||
| - Bootstrap timing: no_peers_retry, per_peer_settle, join_confirmation, discovery_poll | ||
| - Publisher timing: initial_delay, base_interval, max_jitter | ||
| - Merge timing: base_interval, max_jitter (separate for bubble and overlap) | ||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🧩 Analysis chain
🏁 Script executed:
Repository: rustonbsd/distributed-topic-tracker
Length of output: 50393
🏁 Script executed:
Repository: rustonbsd/distributed-topic-tracker
Length of output: 96
🏁 Script executed:
head -n 160 ARCHITECTURE.md | tail -n +100Repository: rustonbsd/distributed-topic-tracker
Length of output: 1717
🏁 Script executed:
sed -n '110,130p' ARCHITECTURE.mdRepository: rustonbsd/distributed-topic-tracker
Length of output: 651
🏁 Script executed:
sed -n '113,125p' ARCHITECTURE.mdRepository: rustonbsd/distributed-topic-tracker
Length of output: 428
🏁 Script executed:
Repository: rustonbsd/distributed-topic-tracker
Length of output: 1344
🏁 Script executed:
sed -n '150,175p' src/gossip/topic/bootstrap.rsRepository: rustonbsd/distributed-topic-tracker
Length of output: 1285
Update pseudocode and documentation to clarify conditional fetch windows based on
check_older_records_first_on_startup.The documented invariant "Both
unix_minuteandunix_minute - 1records are always fetched" contradicts the actual behavior. Whencheck_older_records_first_on_startupis true, the implementation fetchescurrent-2andcurrent-1; when false, it fetchescurrent-1andcurrent. The pseudocode should reflect this conditional logic explicitly, or the documentation should clarify that the fetch window changes based on configuration.🤖 Prompt for AI Agents