Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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
2 changes: 2 additions & 0 deletions .github/workflows/deploy.yml
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,7 @@ jobs:
node test-1659-analytics-warmup.js
node test-app-api-inflight-cleanup-rejection.js
node test-channels-merge-1498-unit.js
node test-channels-observed-path-hash-size.js
node test-issue-1518-home-url.js
node test-channel-decrypt-insecure-context.js
node test-live-region-filter.js
Expand Down Expand Up @@ -564,6 +565,7 @@ jobs:
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-channels-share-color-e2e.js 2>&1 | tee -a e2e-output.txt
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-channels-ws-batch-e2e.js 2>&1 | tee -a e2e-output.txt
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-channels-ws-race-1498-e2e.js 2>&1 | tee -a e2e-output.txt
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-channels-observed-path-hash-size-e2e.js 2>&1 | tee -a e2e-output.txt
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-issue-1487-byop-modal-layout-e2e.js 2>&1 | tee -a e2e-output.txt
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-issue-1630-reach-mobile-e2e.js 2>&1 | tee -a e2e-output.txt
CHROMIUM_REQUIRE=1 BASE_URL=http://localhost:13581 node test-node-reach-coverage-e2e.js 2>&1 | tee -a e2e-output.txt
Expand Down
1 change: 1 addition & 0 deletions cmd/server/chunked_load.go
Original file line number Diff line number Diff line change
Expand Up @@ -637,6 +637,7 @@ func (s *PacketStore) scanAndMergeChunk(rows *sql.Rows, relayPM *prefixMap, cold
}
}

tx.mergeObservedPathHashSize(obsPJ)
tx.Observations = append(tx.Observations, obs)
tx.obsKeys[dk] = true
if obs.ObserverID != "" && !tx.observerSet[obs.ObserverID] {
Expand Down
44 changes: 25 additions & 19 deletions cmd/server/db.go
Original file line number Diff line number Diff line change
Expand Up @@ -3506,9 +3506,10 @@ func (db *DB) GetChannelMessages(channelHash string, limit, offset int, region .
defer rows.Close()

type msg struct {
Data map[string]interface{}
Repeats int
LatestEpoch int64 // max observation timestamp (unix seconds) — issue #1366
Data map[string]interface{}
Repeats int
LatestEpoch int64 // max observation timestamp (unix seconds) — issue #1366
PathHashSizeMask uint8
}
msgMap := make(map[int]*msg, len(pageIDs))

Expand Down Expand Up @@ -3582,8 +3583,11 @@ func (db *DB) GetChannelMessages(channelHash string, limit, offset int, region .
observerName = obsID.String
}

pathHashSizeMask := observedPathHashSizeMask(nullStrVal(pathJSON))
if existing, ok := msgMap[txID]; ok {
existing.Repeats++
existing.PathHashSizeMask |= pathHashSizeMask
existing.Data["observedPathHashSizes"] = observedPathHashSizes(existing.PathHashSizeMask)
if obsTs.Valid && obsTs.Int64 > existing.LatestEpoch {
existing.LatestEpoch = obsTs.Int64
}
Expand Down Expand Up @@ -3638,23 +3642,25 @@ func (db *DB) GetChannelMessages(channelHash string, limit, offset int, region .
}
m := &msg{
Data: map[string]interface{}{
"sender": displaySender,
"text": displayText,
"timestamp": nullStr(fs),
"first_seen": nullStr(fs),
"sender_timestamp": senderTs,
"packetId": pktID,
"packetHash": nullStr(pktHash),
"repeats": 1,
"observers": []string{},
"hops": hops,
"snr": nullFloat(snr),
"scope": nullStr(scopeName),
"routeType": nullInt(routeType),
"entryPrefix": entryPrefix,
"entryObserverPubkey": entryObserverPubkey,
"sender": displaySender,
"text": displayText,
"timestamp": nullStr(fs),
"first_seen": nullStr(fs),
"sender_timestamp": senderTs,
"packetId": pktID,
"packetHash": nullStr(pktHash),
"repeats": 1,
"observers": []string{},
"hops": hops,
"snr": nullFloat(snr),
"scope": nullStr(scopeName),
"routeType": nullInt(routeType),
"entryPrefix": entryPrefix,
"entryObserverPubkey": entryObserverPubkey,
"observedPathHashSizes": observedPathHashSizes(pathHashSizeMask),
},
Repeats: 1,
Repeats: 1,
PathHashSizeMask: pathHashSizeMask,
}
if obsTs.Valid {
m.LatestEpoch = obsTs.Int64
Expand Down
44 changes: 44 additions & 0 deletions cmd/server/hash_migrate.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,17 @@ func migrateContentHashesAsync(store *PacketStore, batchSize int, yieldDuration
total := len(store.packets)
store.mu.RUnlock()

// Keep evidence only for hashes that actually collide during this one
// migration run. The legacy migration retains duplicate in-memory rows, so
// a later batch must also update ghosts created by an earlier collision.
// This registry is bounded by collision members and is discarded when the
// migration returns.
type collisionEvidence struct {
mask uint8
members map[*StoreTx]struct{}
}
collisionEvidenceByHash := make(map[string]*collisionEvidence)

migrated := 0
for offset := 0; offset < total; offset += batchSize {
end := offset + batchSize
Expand Down Expand Up @@ -52,6 +63,12 @@ func migrateContentHashesAsync(store *PacketStore, batchSize int, yieldDuration
if len(updates) == 0 {
continue
}
// A UNIQUE collision merges DB observations into one survivor, while
// this legacy migration intentionally keeps its existing in-memory
// cardinality/index behaviour. Track the survivor IDs so the observed
// path-width evidence can nevertheless be made consistent across every
// same-content in-memory row after the existing update loop.
collisionSurvivors := make(map[string][]int)

// Write batch to DB in a single transaction.
dbTx, err := store.db.conn.Begin()
Expand All @@ -75,6 +92,7 @@ func migrateContentHashesAsync(store *PacketStore, batchSize int, yieldDuration
if err2 := dbTx.QueryRow("SELECT id FROM transmissions WHERE hash = ?", u.newHash).Scan(&survID); err2 == nil {
dbTx.Exec("UPDATE observations SET transmission_id = ? WHERE transmission_id = ?", survID, u.tx.ID)
dbTx.Exec("DELETE FROM transmissions WHERE id = ?", u.tx.ID)
collisionSurvivors[u.newHash] = append(collisionSurvivors[u.newHash], survID)
u.newHash = "" // mark for in-memory removal only
}
}
Expand Down Expand Up @@ -105,6 +123,32 @@ func migrateContentHashesAsync(store *PacketStore, batchSize int, yieldDuration
store.byHash[u.newHash] = u.tx
}
}
for newHash, survivorIDs := range collisionSurvivors {
evidence := collisionEvidenceByHash[newHash]
if evidence == nil {
evidence = &collisionEvidence{members: make(map[*StoreTx]struct{})}
collisionEvidenceByHash[newHash] = evidence
}
addMember := func(tx *StoreTx) {
if tx == nil {
return
}
evidence.mask |= tx.pathHashSizeMask
evidence.members[tx] = struct{}{}
}
for _, u := range updates {
if u.newHash == newHash {
addMember(u.tx)
}
}
for _, survivorID := range survivorIDs {
addMember(store.byTxID[survivorID])
}
addMember(store.byHash[newHash])
for member := range evidence.members {
member.pathHashSizeMask |= evidence.mask
}
}
store.mu.Unlock()

migrated += len(updates)
Expand Down
200 changes: 200 additions & 0 deletions cmd/server/hash_migrate_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package main

import (
"reflect"
"testing"
"time"
)
Expand Down Expand Up @@ -76,3 +77,202 @@ func TestMigrateContentHashesAsync_NoOp(t *testing.T) {
t.Error("hash should remain in index")
}
}

func TestMigrateContentHashesAsync_CollisionUnionsObservedPathHashSizesWithoutChangingLegacyIndexes(t *testing.T) {
db := setupTestDBv2(t)
store := NewPacketStore(db, nil)

rawHex := "0A00D69FD7A5A7475DB07337749AE61FA53A4788E976"
correctHash := ComputeContentHash(rawHex)
for _, row := range []struct {
id int
wrongHash string
firstSeen string
pathJSON string
}{
{1, "old-hash-one", "2026-01-01T00:00:00Z", `["AA"]`},
{2, "old-hash-two", "2026-01-01T00:00:01Z", `["BEEF"]`},
} {
if _, err := db.conn.Exec(`INSERT INTO transmissions
(id, raw_hex, hash, first_seen, route_type, payload_type, decoded_json)
VALUES (?, ?, ?, ?, 1, 5, '{}')`, row.id, rawHex, row.wrongHash, row.firstSeen); err != nil {
t.Fatal(err)
}
if _, err := db.conn.Exec(`INSERT INTO observations
(id, transmission_id, observer_id, observer_name, path_json, timestamp)
VALUES (?, ?, ?, ?, ?, ?)`, row.id, row.id, "observer", "Observer", row.pathJSON, row.id); err != nil {
t.Fatal(err)
}
}

if err := store.Load(); err != nil {
t.Fatal(err)
}
// One row per batch proves that the survivor's evidence is recovered from
// byTxID rather than only from same-batch updates.
migrateContentHashesAsync(store, 1, 0)

wantSizes := []int{1, 2}
for _, tx := range store.packets {
if got := tx.observedPathHashSizes(); !reflect.DeepEqual(got, wantSizes) {
t.Fatalf("tx %d observed path hash sizes = %v, want %v", tx.ID, got, wantSizes)
}
}
if authoritative := store.byHash[correctHash]; authoritative == nil ||
!reflect.DeepEqual(authoritative.observedPathHashSizes(), wantSizes) {
t.Fatalf("authoritative observed path hash sizes = %v, want %v",
authoritative.observedPathHashSizes(), wantSizes)
}

// Characterize, do not silently fix, the pre-existing duplicate-migration
// index bug. A separate change must remove the deleted row from every
// PacketStore index and re-parent its observations atomically.
if len(store.packets) != 2 || len(store.byPayloadType[PayloadGRP_TXT]) != 2 {
t.Fatalf("legacy in-memory cardinality changed: packets/payload = %d/%d, want 2/2",
len(store.packets), len(store.byPayloadType[PayloadGRP_TXT]))
}
if got := store.QueryPackets(PacketQuery{Limit: 10}).Total; got != 2 {
t.Fatalf("legacy QueryPackets cardinality changed: got %d, want 2", got)
}
if store.byTxID[2] == nil || store.byTxID[2].Observations[0].TransmissionID != 2 {
t.Fatal("legacy duplicate/observation ownership changed unexpectedly")
}

var txCount, canonicalObservationCount int
if err := db.conn.QueryRow(`SELECT COUNT(*) FROM transmissions WHERE raw_hex = ?`, rawHex).Scan(&txCount); err != nil {
t.Fatal(err)
}
if txCount != 1 {
t.Fatalf("DB transmission rows = %d, want 1 after duplicate merge", txCount)
}
if err := db.conn.QueryRow(`SELECT COUNT(*) FROM observations WHERE transmission_id = 1`).Scan(&canonicalObservationCount); err != nil {
t.Fatal(err)
}
if canonicalObservationCount != 2 {
t.Fatalf("DB canonical observations = %d, want 2 after duplicate merge", canonicalObservationCount)
}
}

func TestMigrateContentHashesAsync_CollisionIncludesAlreadyCurrentSurvivorEvidence(t *testing.T) {
db := setupTestDBv2(t)
store := NewPacketStore(db, nil)

rawHex := "0A00D69FD7A5A7475DB07337749AE61FA53A4788E976"
correctHash := ComputeContentHash(rawHex)
for _, row := range []struct {
id int
hash string
pathJSON string
}{
{1, correctHash, `["010203"]`},
{2, "old-hash", `["BEEF"]`},
} {
if _, err := db.conn.Exec(`INSERT INTO transmissions
(id, raw_hex, hash, first_seen, route_type, payload_type, decoded_json)
VALUES (?, ?, ?, ?, 1, 5, '{}')`, row.id, rawHex, row.hash,
time.Date(2026, 1, 1, 0, 0, row.id, 0, time.UTC).Format(time.RFC3339)); err != nil {
t.Fatal(err)
}
if _, err := db.conn.Exec(`INSERT INTO observations
(id, transmission_id, observer_id, observer_name, path_json, timestamp)
VALUES (?, ?, ?, ?, ?, ?)`, row.id, row.id, "observer", "Observer", row.pathJSON, row.id); err != nil {
t.Fatal(err)
}
}

if err := store.Load(); err != nil {
t.Fatal(err)
}
migrateContentHashesAsync(store, 100, 0)

want := []int{2, 3}
for _, tx := range store.packets {
if got := tx.observedPathHashSizes(); !reflect.DeepEqual(got, want) {
t.Fatalf("tx %d observed path hash sizes = %v, want %v", tx.ID, got, want)
}
}
if got := store.byHash[correctHash].observedPathHashSizes(); !reflect.DeepEqual(got, want) {
t.Fatalf("authoritative observed path hash sizes = %v, want %v", got, want)
}
if len(store.packets) != 2 || store.QueryPackets(PacketQuery{Limit: 10}).Total != 2 {
t.Fatal("legacy in-memory duplicate cardinality changed unexpectedly")
}
}

func TestMigrateContentHashesAsync_CollisionCarriesEvidenceAcrossBatches(t *testing.T) {
db := setupTestDBv2(t)
store := NewPacketStore(db, nil)

rawHex := "0A00D69FD7A5A7475DB07337749AE61FA53A4788E976"
correctHash := ComputeContentHash(rawHex)
for _, row := range []struct {
id int
wrongHash string
pathJSON string
}{
{1, "old-hash-one", `["AA"]`},
{2, "old-hash-two", `["BEEF"]`},
{3, "old-hash-three", `["010203"]`},
} {
if _, err := db.conn.Exec(`INSERT INTO transmissions
(id, raw_hex, hash, first_seen, route_type, payload_type, decoded_json)
VALUES (?, ?, ?, ?, 1, 5, '{}')`, row.id, rawHex, row.wrongHash,
time.Date(2026, 1, 1, 0, 0, row.id, 0, time.UTC).Format(time.RFC3339)); err != nil {
t.Fatal(err)
}
if _, err := db.conn.Exec(`INSERT INTO observations
(id, transmission_id, observer_id, observer_name, path_json, timestamp)
VALUES (?, ?, ?, ?, ?, ?)`, row.id, row.id, "observer", "Observer", row.pathJSON, row.id); err != nil {
t.Fatal(err)
}
}

if err := store.Load(); err != nil {
t.Fatal(err)
}
// Force each collision into a different transaction/batch. Evidence learned
// by the third collision must flow back into the ghost row retained by the
// second collision's legacy in-memory behaviour.
migrateContentHashesAsync(store, 1, 0)

want := []int{1, 2, 3}
for _, tx := range store.packets {
if got := tx.observedPathHashSizes(); !reflect.DeepEqual(got, want) {
t.Fatalf("tx %d observed path hash sizes = %v, want %v", tx.ID, got, want)
}
}
if got := store.byHash[correctHash].observedPathHashSizes(); !reflect.DeepEqual(got, want) {
t.Fatalf("authoritative observed path hash sizes = %v, want %v", got, want)
}

// Keep characterizing the separate legacy index bug rather than hiding it
// inside this evidence-only feature change.
if len(store.packets) != 3 || len(store.byPayloadType[PayloadGRP_TXT]) != 3 {
t.Fatalf("legacy in-memory cardinality changed: packets/payload = %d/%d, want 3/3",
len(store.packets), len(store.byPayloadType[PayloadGRP_TXT]))
}
if got := store.QueryPackets(PacketQuery{Limit: 10}).Total; got != 3 {
t.Fatalf("legacy QueryPackets cardinality changed: got %d, want 3", got)
}
for _, duplicateID := range []int{2, 3} {
duplicate := store.byTxID[duplicateID]
if duplicate == nil || len(duplicate.Observations) != 1 ||
duplicate.Observations[0].TransmissionID != duplicateID {
t.Fatalf("legacy duplicate %d observation ownership changed unexpectedly", duplicateID)
}
}

var txCount, canonicalObservationCount int
if err := db.conn.QueryRow(`SELECT COUNT(*) FROM transmissions WHERE raw_hex = ?`, rawHex).Scan(&txCount); err != nil {
t.Fatal(err)
}
if txCount != 1 {
t.Fatalf("DB transmission rows = %d, want 1 after duplicate merges", txCount)
}
if err := db.conn.QueryRow(`SELECT COUNT(*) FROM observations WHERE transmission_id = 1`).Scan(&canonicalObservationCount); err != nil {
t.Fatal(err)
}
if canonicalObservationCount != 3 {
t.Fatalf("DB canonical observations = %d, want 3 after duplicate merges", canonicalObservationCount)
}
}
Loading
Loading