Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
eabc354
chore: dep update actor-helper
rustonbsd Apr 12, 2026
a6a6a79
fix: cleanup
rustonbsd Apr 12, 2026
b58fdfe
fix: cicd actions/checkout update to v6
rustonbsd Apr 12, 2026
391abcc
fix: on drop of topic, sender and receiver pointer stops all backgrou…
rustonbsd Apr 13, 2026
a0c73fa
chore: explicit handeling of action error in actor run loops + fixed…
rustonbsd Apr 13, 2026
93b909e
chore: added timeouts to sender calls that can halt due to full iroh-…
rustonbsd Apr 14, 2026
1de3c6e
fix: redundent test removed from cicd and removed edited from on pull…
rustonbsd Apr 14, 2026
8afe88f
chore: test_topic_full_shutdown_on_drop test using cancel_token.cance…
rustonbsd Apr 14, 2026
10303f4
chore: receiver functions converted to Result<>
rustonbsd Apr 14, 2026
bcb94c6
add: config everything + fixes all around
rustonbsd Apr 16, 2026
7b210c3
fix: cargo test_gossip removed
rustonbsd Apr 16, 2026
341ff8c
fix: next_interval across project min(1000) sec to prevent hot loops …
rustonbsd Apr 16, 2026
24c8bb5
fix: review fixes
rustonbsd Apr 16, 2026
66158a5
fix: explicitly set bootstrap first record unix_minute with config
rustonbsd Apr 16, 2026
44163f8
fix: next and join sender reference counting + changed base_interval …
rustonbsd Apr 16, 2026
a9d0ce1
chore: review cleanup
rustonbsd Apr 16, 2026
4d7bad7
chore: impl review suggestions + fix and add docs
rustonbsd Apr 17, 2026
1ba92ff
chore: more review fixes
rustonbsd Apr 17, 2026
68557ad
chore: more review stuff + one less dht get call on startup (cached r…
rustonbsd Apr 18, 2026
44f2697
fix: dht put seq_num + topic::new async_bootstrap + error passthrough…
rustonbsd Apr 18, 2026
2d6d793
chore: messave_overlap changed from miss on single message to miss on…
rustonbsd Apr 18, 2026
07405ac
chore: update docs and correctness fixes
rustonbsd Apr 18, 2026
b410eca
fix: topic error propagation + topic fail_on_* now fully enforced + c…
rustonbsd Apr 19, 2026
fa81706
chore: docs cleanup + wait for spawn_workers tasks
rustonbsd Apr 19, 2026
7701583
chore: review fixes
rustonbsd Apr 19, 2026
228de6e
fix: moved TopicId to crypto + fixed interval jitter calc + check_las…
rustonbsd Apr 19, 2026
630e8a1
chore: review fixes
rustonbsd Apr 19, 2026
6d82584
add: cancelation token to get/put record + minor review docs fixes
rustonbsd Apr 19, 2026
a928778
chore: streamlined > 0 enforcement, removed ZeroU32
rustonbsd Apr 19, 2026
e29f18a
chore: interval jitter calculation now uses native u128
rustonbsd Apr 19, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions .github/workflows/test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ on:

# always on pullrequest
pull_request:
types: [opened, synchronize, reopened, edited]
types: [opened, synchronize, reopened]


env:
Expand All @@ -18,7 +18,7 @@ jobs:
runs-on: ubuntu-latest

steps:
- uses: actions/checkout@v4
- uses: actions/checkout@v6

- name: Build Docker image
run: docker build -t distributed-topic-tracker .
Expand Down
127 changes: 73 additions & 54 deletions ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand All @@ -35,6 +37,8 @@ flowchart LR
subgraph AutoDiscovery
B[Bootstrap Loop]
P[Publisher]
BM[Bubble Merge]
MO[Message Overlap Merge]
end

subgraph DHT
Expand All @@ -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:

Expand All @@ -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
```
Expand All @@ -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:
Expand All @@ -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))
Comment on lines +123 to +124

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

🧩 Analysis chain

🏁 Script executed:

#!/bin/bash
rg -nP --type=rust -C5 'check_older_records_first_on_startup|unix_minute'

Repository: rustonbsd/distributed-topic-tracker

Length of output: 50393


🏁 Script executed:

find . -name "ARCHITECTURE.md" -o -name "architecture.md" | head -20

Repository: rustonbsd/distributed-topic-tracker

Length of output: 96


🏁 Script executed:

head -n 160 ARCHITECTURE.md | tail -n +100

Repository: rustonbsd/distributed-topic-tracker

Length of output: 1717


🏁 Script executed:

sed -n '110,130p' ARCHITECTURE.md

Repository: rustonbsd/distributed-topic-tracker

Length of output: 651


🏁 Script executed:

sed -n '113,125p' ARCHITECTURE.md

Repository: rustonbsd/distributed-topic-tracker

Length of output: 428


🏁 Script executed:

sed -n '1,130p' ARCHITECTURE.md | tail -n 40

Repository: rustonbsd/distributed-topic-tracker

Length of output: 1344


🏁 Script executed:

sed -n '150,175p' src/gossip/topic/bootstrap.rs

Repository: 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_minute and unix_minute - 1 records are always fetched" contradicts the actual behavior. When check_older_records_first_on_startup is true, the implementation fetches current-2 and current-1; when false, it fetches current-1 and current. 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
Verify each finding against the current code and only fix it if needed.

In `@ARCHITECTURE.md` around lines 123 - 124, The pseudocode and docs currently
state that both unix_minute and unix_minute - 1 are always fetched but the
implementation uses a conditional window; update the pseudocode around the
minute calculation (variables minute, first_attempt) and the conditional flag
(check_older_records_first_on_startup / check_last_minute_first) so it
explicitly shows both cases: when check_older_records_first_on_startup is true
fetch get_records(unix_minute(minute) - 2) and get_records(unix_minute(minute) -
1) (i.e., current-2 and current-1), and when false fetch
get_records(unix_minute(minute) - 1) and get_records(unix_minute(minute)) (i.e.,
current-1 and current); ensure unix_minute and get_records are referenced
exactly and update the invariant text to match this conditional behavior.


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
Expand All @@ -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).

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟡 Minor

Documented parameter name doesn't match the code.

The docs call the bubble-merge join cap max_join_peer_count, but the actual field/parameter in src/gossip/merge/bubble.rs is max_join_peers (see the BubbleMerge::new signature and BubbleMergeActor struct). Readers trying to map the ARCHITECTURE.md parameters to Config/constructor fields will get confused. Please unify the name.

📝 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
Verify each finding against the current code and only fix it if needed.

In `@ARCHITECTURE.md` at line 179, The docs use the parameter name
max_join_peer_count but the code (BubbleMerge::new and the BubbleMergeActor
struct) uses max_join_peers; update them to match by choosing one canonical name
and making both the ARCHITECTURE.md entry (and any other doc occurrences) and
the code signatures/field names consistent — e.g., rename the doc entry to
max_join_peers or refactor the code to max_join_peer_count and update
BubbleMerge::new and the BubbleMergeActor field accordingly, plus any tests or
usages that reference the old name.


Signal 2: non-overlapping message sets.
- Compare local last_message_hashes with others.
Expand All @@ -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
```

Expand All @@ -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:
Expand All @@ -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:
Expand All @@ -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
Expand All @@ -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)
Loading