fix(gossip): persist USER_UPSERT rows and use shared RPC broadcast helper
Fill required user schema fields on ingest and route mesh gossip through broadcastRpcToSessions. Co-authored-by: Cursor <[email protected]>
This commit is contained in:
@@ -253,8 +253,16 @@ class PearcordGuild extends EventEmitter {
|
||||
this.emit('presence', { peerId, ...payload })
|
||||
}
|
||||
if (method === RPC.USER_UPSERT && payload?.id) {
|
||||
this.db.insert(COLLECTIONS.USERS, payload).then(() => {
|
||||
this.emit('user', payload)
|
||||
const row = {
|
||||
discriminator: '0000',
|
||||
publicKey: '',
|
||||
createdAt: now(),
|
||||
...payload,
|
||||
id: payload.id
|
||||
}
|
||||
if (!row.username) row.username = payload.id.slice(0, 8)
|
||||
this.db.insert(COLLECTIONS.USERS, row).then(() => {
|
||||
this.emit('user', row)
|
||||
}).catch(() => {})
|
||||
}
|
||||
if (method === RPC.PROFILE_COSMETIC_UPDATE) {
|
||||
|
||||
@@ -3,8 +3,8 @@
|
||||
const Protomux = require('protomux')
|
||||
const b4a = require('b4a')
|
||||
const c = require('compact-encoding')
|
||||
const { RPC, encodeRpc, decodeRpc, wireSwarmConnection } = require('pearcord-shared')
|
||||
const { openWireChannel, wireReady } = require('pearcord-drive/mux-wire')
|
||||
const { decodeRpc, wireSwarmConnection } = require('pearcord-shared')
|
||||
const { openWireChannel, broadcastRpcToSessions } = require('pearcord-drive/mux-wire')
|
||||
|
||||
const GOSSIP_PROTOCOL = 'pearcord-gossip-v1'
|
||||
|
||||
@@ -29,25 +29,13 @@ function attachGossipMesh (guild, conn) {
|
||||
}
|
||||
|
||||
async function broadcastGossip (guild, method, payload) {
|
||||
const m = Number(method)
|
||||
if (!Number.isInteger(m) || m <= 0 || m > 255) return
|
||||
let buf
|
||||
try {
|
||||
buf = encodeRpc(m, payload)
|
||||
} catch {
|
||||
return
|
||||
}
|
||||
const sessions = [...guild._channels.values()]
|
||||
await Promise.all(
|
||||
sessions.map(async (session) => {
|
||||
try {
|
||||
await wireReady(session, 2000)
|
||||
session.send(buf)
|
||||
} catch {
|
||||
// peer gone
|
||||
}
|
||||
})
|
||||
)
|
||||
await broadcastRpcToSessions(guild._channels, method, payload, { wireReadyMs: 2000 })
|
||||
}
|
||||
|
||||
module.exports = { attachGossipMesh, broadcastGossip, GOSSIP_PROTOCOL, rpcPayload }
|
||||
module.exports = {
|
||||
attachGossipMesh,
|
||||
broadcastGossip,
|
||||
broadcastRpcToSessions,
|
||||
GOSSIP_PROTOCOL,
|
||||
rpcPayload
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user