Fix connection teardown and sync lifecycle on login
This commit is contained in:
parent
ae5992c948
commit
eae644ab3c
3 changed files with 41 additions and 12 deletions
|
|
@ -147,9 +147,10 @@ export const makeFeed = ({
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
|
seen.add(event.id)
|
||||||
|
|
||||||
if (between([backwardWindow[0], forwardWindow[1]], event.created_at)) {
|
if (between([backwardWindow[0], forwardWindow[1]], event.created_at)) {
|
||||||
visible.push(event)
|
visible.push(event)
|
||||||
seen.add(event.id)
|
|
||||||
} else {
|
} else {
|
||||||
insertIntoBuffer(event)
|
insertIntoBuffer(event)
|
||||||
}
|
}
|
||||||
|
|
@ -176,6 +177,18 @@ export const makeFeed = ({
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Buffered events are routed through insertEvents again, so forget we've seen
|
||||||
|
// them to let the window check run a second time
|
||||||
|
const drainBuffer = () => {
|
||||||
|
const drained = buffer.splice(0, 30)
|
||||||
|
|
||||||
|
for (const event of drained) {
|
||||||
|
seen.delete(event.id)
|
||||||
|
}
|
||||||
|
|
||||||
|
insertEvents(drained)
|
||||||
|
}
|
||||||
|
|
||||||
const unsubscribers = [
|
const unsubscribers = [
|
||||||
on(
|
on(
|
||||||
repository,
|
repository,
|
||||||
|
|
@ -239,7 +252,7 @@ export const makeFeed = ({
|
||||||
|
|
||||||
backwardWindow = [since - interval, since]
|
backwardWindow = [since - interval, since]
|
||||||
|
|
||||||
insertEvents(buffer.splice(0, 30))
|
drainBuffer()
|
||||||
|
|
||||||
if (until > now() - int(2, YEAR)) {
|
if (until > now() - int(2, YEAR)) {
|
||||||
await loadTimeframe(since, until)
|
await loadTimeframe(since, until)
|
||||||
|
|
@ -260,7 +273,7 @@ export const makeFeed = ({
|
||||||
|
|
||||||
forwardWindow = [until, until + interval]
|
forwardWindow = [until, until + interval]
|
||||||
|
|
||||||
insertEvents(buffer.splice(0, 30))
|
drainBuffer()
|
||||||
|
|
||||||
if (until < now()) {
|
if (until < now()) {
|
||||||
await loadTimeframe(since, until)
|
await loadTimeframe(since, until)
|
||||||
|
|
|
||||||
|
|
@ -1,6 +1,6 @@
|
||||||
import {page} from "$app/stores"
|
import {page} from "$app/stores"
|
||||||
import type {Unsubscriber} from "svelte/store"
|
import type {Unsubscriber} from "svelte/store"
|
||||||
import {call, assoc, WEEK, MONTH, ago} from "@welshman/lib"
|
import {last, call, assoc, WEEK, MONTH, ago} from "@welshman/lib"
|
||||||
import {merged} from "@welshman/store"
|
import {merged} from "@welshman/store"
|
||||||
import {Router} from "@welshman/router"
|
import {Router} from "@welshman/router"
|
||||||
import {
|
import {
|
||||||
|
|
@ -82,6 +82,7 @@ const pullOneWithFallback = async (
|
||||||
if (signal.aborted) return
|
if (signal.aborted) return
|
||||||
|
|
||||||
const cachedEvents = repository.query([filter]).filter(isSignedEvent)
|
const cachedEvents = repository.query([filter]).filter(isSignedEvent)
|
||||||
|
const since = last(cachedEvents.slice(10))?.created_at || 0
|
||||||
|
|
||||||
if (onEvent) {
|
if (onEvent) {
|
||||||
for (const event of cachedEvents) {
|
for (const event of cachedEvents) {
|
||||||
|
|
@ -116,7 +117,7 @@ const pullOneWithFallback = async (
|
||||||
// }))
|
// }))
|
||||||
|
|
||||||
if (shouldFallback && !signal.aborted) {
|
if (shouldFallback && !signal.aborted) {
|
||||||
request({relays: [url], signal, autoClose: true, filters: [filter], onEvent})
|
request({relays: [url], signal, autoClose: true, filters: [{since, ...filter}], onEvent})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -468,17 +469,20 @@ const syncDMs = () => {
|
||||||
|
|
||||||
// Merge all synchronization functions
|
// Merge all synchronization functions
|
||||||
|
|
||||||
let unsubscribe: Unsubscriber
|
let unsubscribe: Unsubscriber | undefined
|
||||||
|
|
||||||
export const syncApplicationData = () => {
|
export const syncApplicationData = () => {
|
||||||
const unsubscribers = [syncRelays(), syncUserData(), syncSpaces(), syncDMs()]
|
const unsubscribers = [syncRelays(), syncUserData(), syncSpaces(), syncDMs()]
|
||||||
|
|
||||||
unsubscribe = () => unsubscribers.forEach(call)
|
unsubscribe = () => unsubscribers.forEach(call)
|
||||||
|
}
|
||||||
|
|
||||||
return unsubscribe
|
export const stopApplicationDataSync = () => {
|
||||||
|
unsubscribe?.()
|
||||||
|
unsubscribe = undefined
|
||||||
}
|
}
|
||||||
|
|
||||||
export const resyncApplicationData = () => {
|
export const resyncApplicationData = () => {
|
||||||
unsubscribe?.()
|
stopApplicationDataSync()
|
||||||
syncApplicationData()
|
syncApplicationData()
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -38,7 +38,7 @@
|
||||||
import {getSetting, userSettings, notificationSettings} from "@app/settings"
|
import {getSetting, userSettings, notificationSettings} from "@app/settings"
|
||||||
import {DUFFLEPUD_URL, DEFAULT_RELAYS, INDEXER_RELAYS, POMADE_SIGNERS} from "@app/env"
|
import {DUFFLEPUD_URL, DEFAULT_RELAYS, INDEXER_RELAYS, POMADE_SIGNERS} from "@app/env"
|
||||||
import {pushState} from "@app/push/adapters/common"
|
import {pushState} from "@app/push/adapters/common"
|
||||||
import {syncApplicationData} from "@app/sync"
|
import {syncApplicationData, resyncApplicationData, stopApplicationDataSync} from "@app/sync"
|
||||||
import * as groups from "@app/groups"
|
import * as groups from "@app/groups"
|
||||||
import * as comments from "@app/comments"
|
import * as comments from "@app/comments"
|
||||||
import * as deletes from "@app/deletes"
|
import * as deletes from "@app/deletes"
|
||||||
|
|
@ -250,7 +250,9 @@
|
||||||
unsubscribers.push(() => defaultSocketPolicies.splice(-policies.length))
|
unsubscribers.push(() => defaultSocketPolicies.splice(-policies.length))
|
||||||
|
|
||||||
// History, navigation, application data
|
// History, navigation, application data
|
||||||
unsubscribers.push(setupHistory(), setupAnalytics(), syncApplicationData())
|
syncApplicationData()
|
||||||
|
|
||||||
|
unsubscribers.push(setupHistory(), setupAnalytics(), stopApplicationDataSync)
|
||||||
|
|
||||||
// Initialize keyboard state tracking
|
// Initialize keyboard state tracking
|
||||||
unsubscribers.push(syncKeyboard())
|
unsubscribers.push(syncKeyboard())
|
||||||
|
|
@ -267,8 +269,18 @@
|
||||||
// Initialize background notifications
|
// Initialize background notifications
|
||||||
unsubscribers.push(Push.sync())
|
unsubscribers.push(Push.sync())
|
||||||
|
|
||||||
// Any time our pubkey changes, close all connections
|
// When the user logs in, drop connections opened anonymously and sync again
|
||||||
pubkey.subscribe(() => Pool.get().clear())
|
let lastPubkey = pubkey.get()
|
||||||
|
|
||||||
|
unsubscribers.push(
|
||||||
|
pubkey.subscribe($pubkey => {
|
||||||
|
if ($pubkey !== lastPubkey) {
|
||||||
|
lastPubkey = $pubkey
|
||||||
|
Pool.get().clear()
|
||||||
|
resyncApplicationData()
|
||||||
|
}
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
|
||||||
// Listen for signer errors, report to user via toast
|
// Listen for signer errors, report to user via toast
|
||||||
unsubscribers.push(
|
unsubscribers.push(
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue