Handle init/resume/suspend/background and polling
This commit is contained in:
+135
-97
@@ -190,7 +190,6 @@ export class Convo {
|
|||||||
private historyCursor: string | undefined | null = undefined
|
private historyCursor: string | undefined | null = undefined
|
||||||
private isFetchingHistory = false
|
private isFetchingHistory = false
|
||||||
private eventsCursor: string | undefined = undefined
|
private eventsCursor: string | undefined = undefined
|
||||||
private pollingFailure = false
|
|
||||||
|
|
||||||
private pastMessages: Map<
|
private pastMessages: Map<
|
||||||
string,
|
string,
|
||||||
@@ -208,8 +207,9 @@ export class Convo {
|
|||||||
private footerItems: Map<string, ConvoItem> = new Map()
|
private footerItems: Map<string, ConvoItem> = new Map()
|
||||||
private headerItems: Map<string, ConvoItem> = new Map()
|
private headerItems: Map<string, ConvoItem> = new Map()
|
||||||
|
|
||||||
private pendingEventIngestion: Promise<void> | undefined
|
|
||||||
private isProcessingPendingMessages = false
|
private isProcessingPendingMessages = false
|
||||||
|
private pendingPoll: Promise<void> | undefined
|
||||||
|
private nextPoll: NodeJS.Timeout | undefined
|
||||||
|
|
||||||
convoId: string
|
convoId: string
|
||||||
convo: ChatBskyConvoDefs.ConvoView | undefined
|
convo: ChatBskyConvoDefs.ConvoView | undefined
|
||||||
@@ -217,7 +217,10 @@ export class Convo {
|
|||||||
recipients: AppBskyActorDefs.ProfileViewBasic[] | undefined = undefined
|
recipients: AppBskyActorDefs.ProfileViewBasic[] | undefined = undefined
|
||||||
snapshot: ConvoState | undefined
|
snapshot: ConvoState | undefined
|
||||||
|
|
||||||
|
id: string
|
||||||
|
|
||||||
constructor(params: ConvoParams) {
|
constructor(params: ConvoParams) {
|
||||||
|
this.id = nanoid(2)
|
||||||
this.convoId = params.convoId
|
this.convoId = params.convoId
|
||||||
this.agent = params.agent
|
this.agent = params.agent
|
||||||
this.__tempFromUserDid = params.__tempFromUserDid
|
this.__tempFromUserDid = params.__tempFromUserDid
|
||||||
@@ -227,6 +230,12 @@ export class Convo {
|
|||||||
this.sendMessage = this.sendMessage.bind(this)
|
this.sendMessage = this.sendMessage.bind(this)
|
||||||
this.deleteMessage = this.deleteMessage.bind(this)
|
this.deleteMessage = this.deleteMessage.bind(this)
|
||||||
this.fetchMessageHistory = this.fetchMessageHistory.bind(this)
|
this.fetchMessageHistory = this.fetchMessageHistory.bind(this)
|
||||||
|
|
||||||
|
logger.debug(
|
||||||
|
'Convo: created',
|
||||||
|
{convoId: this.convoId, id: this.id},
|
||||||
|
logger.DebugContext.convo,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
private commit() {
|
private commit() {
|
||||||
@@ -318,12 +327,12 @@ export class Convo {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async init() {
|
async init() {
|
||||||
logger.debug('Convo: init', {}, logger.DebugContext.convo)
|
|
||||||
|
|
||||||
if (
|
if (
|
||||||
this.status === ConvoStatus.Uninitialized ||
|
this.status === ConvoStatus.Uninitialized ||
|
||||||
this.status === ConvoStatus.Error
|
this.status === ConvoStatus.Error
|
||||||
) {
|
) {
|
||||||
|
logger.debug('Convo: init', {}, logger.DebugContext.convo)
|
||||||
|
|
||||||
try {
|
try {
|
||||||
this.status = ConvoStatus.Initializing
|
this.status = ConvoStatus.Initializing
|
||||||
this.commit()
|
this.commit()
|
||||||
@@ -348,21 +357,25 @@ export class Convo {
|
|||||||
this.status = ConvoStatus.Error
|
this.status = ConvoStatus.Error
|
||||||
this.commit()
|
this.commit()
|
||||||
}
|
}
|
||||||
|
} else if (this.status === ConvoStatus.Suspended) {
|
||||||
|
this.resume()
|
||||||
} else {
|
} else {
|
||||||
logger.warn(`Convo: cannot init from ${this.status}`)
|
logger.warn(`Convo: cannot init from ${this.status}`)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async resume() {
|
async resume() {
|
||||||
logger.debug('Convo: resume', {}, logger.DebugContext.convo)
|
|
||||||
|
|
||||||
if (
|
if (
|
||||||
this.status === ConvoStatus.Suspended ||
|
this.status === ConvoStatus.Suspended ||
|
||||||
this.status === ConvoStatus.Backgrounded
|
this.status === ConvoStatus.Backgrounded
|
||||||
) {
|
) {
|
||||||
|
logger.debug('Convo: resume', {}, logger.DebugContext.convo)
|
||||||
|
|
||||||
const fromStatus = this.status
|
const fromStatus = this.status
|
||||||
|
|
||||||
try {
|
try {
|
||||||
|
this.cancelNextPoll()
|
||||||
|
|
||||||
this.status = ConvoStatus.Resuming
|
this.status = ConvoStatus.Resuming
|
||||||
this.commit()
|
this.commit()
|
||||||
|
|
||||||
@@ -391,21 +404,36 @@ export class Convo {
|
|||||||
this.commit()
|
this.commit()
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
logger.warn(`Convo: cannot resume from ${this.status}`)
|
logger.debug(
|
||||||
|
`Convo: resume called from ${this.status}`,
|
||||||
|
{},
|
||||||
|
logger.DebugContext.convo,
|
||||||
|
)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async background() {
|
async background() {
|
||||||
logger.debug('Convo: backgrounded', {}, logger.DebugContext.convo)
|
if (
|
||||||
this.status = ConvoStatus.Backgrounded
|
this.status === ConvoStatus.Ready ||
|
||||||
this.pollInterval = BACKGROUND_POLL_INTERVAL
|
this.status === ConvoStatus.Resuming
|
||||||
this.commit()
|
) {
|
||||||
|
logger.debug('Convo: backgrounded', {}, logger.DebugContext.convo)
|
||||||
|
this.status = ConvoStatus.Backgrounded
|
||||||
|
this.pollInterval = BACKGROUND_POLL_INTERVAL
|
||||||
|
this.commit()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async suspend() {
|
async suspend() {
|
||||||
logger.debug('Convo: suspended', {}, logger.DebugContext.convo)
|
if (
|
||||||
this.status = ConvoStatus.Suspended
|
this.status === ConvoStatus.Ready ||
|
||||||
this.commit()
|
this.status === ConvoStatus.Backgrounded ||
|
||||||
|
this.status === ConvoStatus.Resuming
|
||||||
|
) {
|
||||||
|
logger.debug('Convo: suspended', {}, logger.DebugContext.convo)
|
||||||
|
this.status = ConvoStatus.Suspended
|
||||||
|
this.commit()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async refreshConvo() {
|
async refreshConvo() {
|
||||||
@@ -518,115 +546,125 @@ export class Convo {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private async pollEvents() {
|
private async pollEvents() {
|
||||||
|
if (this.pendingPoll) return
|
||||||
|
|
||||||
if (
|
if (
|
||||||
this.status === ConvoStatus.Ready ||
|
this.status === ConvoStatus.Ready ||
|
||||||
this.status === ConvoStatus.Backgrounded
|
this.status === ConvoStatus.Backgrounded
|
||||||
) {
|
) {
|
||||||
if (this.pendingEventIngestion) return
|
logger.debug(
|
||||||
|
'Convo: poll events',
|
||||||
|
{convoId: this.convoId, id: this.id},
|
||||||
|
logger.DebugContext.convo,
|
||||||
|
)
|
||||||
|
|
||||||
/*
|
try {
|
||||||
* Represents a failed state, which is retryable.
|
this.pendingPoll = this.ingestLatestEvents()
|
||||||
*/
|
await this.pendingPoll
|
||||||
if (this.pollingFailure) return
|
this.pendingPoll = undefined
|
||||||
|
this.nextPoll = setTimeout(() => {
|
||||||
|
this.pollEvents()
|
||||||
|
}, this.pollInterval)
|
||||||
|
} catch (e: any) {
|
||||||
|
logger.error('Convo: poll events failed')
|
||||||
|
|
||||||
setTimeout(async () => {
|
this.cancelNextPoll()
|
||||||
this.pendingEventIngestion = this.ingestLatestEvents()
|
this.pendingPoll = undefined
|
||||||
await this.pendingEventIngestion
|
|
||||||
this.pendingEventIngestion = undefined
|
this.footerItems.set(ConvoItemError.PollFailed, {
|
||||||
this.pollEvents()
|
type: 'error-recoverable',
|
||||||
}, this.pollInterval)
|
key: ConvoItemError.PollFailed,
|
||||||
|
code: ConvoItemError.PollFailed,
|
||||||
|
retry: () => {
|
||||||
|
this.footerItems.delete(ConvoItemError.PollFailed)
|
||||||
|
this.commit()
|
||||||
|
this.pollEvents()
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
this.commit()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
private cancelNextPoll() {
|
||||||
|
if (this.nextPoll) clearTimeout(this.nextPoll)
|
||||||
}
|
}
|
||||||
|
|
||||||
async ingestLatestEvents() {
|
async ingestLatestEvents() {
|
||||||
try {
|
// throw new Error('UNCOMMENT TO TEST POLL FAILURE')
|
||||||
// throw new Error('UNCOMMENT TO TEST POLL FAILURE')
|
const response = await this.agent.api.chat.bsky.convo.getLog(
|
||||||
const response = await this.agent.api.chat.bsky.convo.getLog(
|
{
|
||||||
{
|
cursor: this.eventsCursor,
|
||||||
cursor: this.eventsCursor,
|
},
|
||||||
|
{
|
||||||
|
headers: {
|
||||||
|
Authorization: this.__tempFromUserDid,
|
||||||
},
|
},
|
||||||
{
|
},
|
||||||
headers: {
|
)
|
||||||
Authorization: this.__tempFromUserDid,
|
const {logs} = response.data
|
||||||
},
|
|
||||||
},
|
|
||||||
)
|
|
||||||
const {logs} = response.data
|
|
||||||
|
|
||||||
let needsCommit = false
|
let needsCommit = false
|
||||||
|
|
||||||
for (const log of logs) {
|
for (const log of logs) {
|
||||||
|
/*
|
||||||
|
* If there's a rev, we should handle it. If there's not a rev, we don't
|
||||||
|
* know what it is.
|
||||||
|
*/
|
||||||
|
if (typeof log.rev === 'string') {
|
||||||
/*
|
/*
|
||||||
* If there's a rev, we should handle it. If there's not a rev, we don't
|
* We only care about new events
|
||||||
* know what it is.
|
|
||||||
*/
|
*/
|
||||||
if (typeof log.rev === 'string') {
|
if (log.rev > (this.eventsCursor = this.eventsCursor || log.rev)) {
|
||||||
/*
|
/*
|
||||||
* We only care about new events
|
* Update rev regardless of if it's a log type we care about or not
|
||||||
*/
|
*/
|
||||||
if (log.rev > (this.eventsCursor = this.eventsCursor || log.rev)) {
|
this.eventsCursor = log.rev
|
||||||
/*
|
|
||||||
* Update rev regardless of if it's a log type we care about or not
|
|
||||||
*/
|
|
||||||
this.eventsCursor = log.rev
|
|
||||||
|
|
||||||
/*
|
/*
|
||||||
* This is VERY important. We don't want to insert any messages from
|
* This is VERY important. We don't want to insert any messages from
|
||||||
* your other chats.
|
* your other chats.
|
||||||
*/
|
*/
|
||||||
if (log.convoId !== this.convoId) continue
|
if (log.convoId !== this.convoId) continue
|
||||||
|
|
||||||
if (
|
if (
|
||||||
ChatBskyConvoDefs.isLogCreateMessage(log) &&
|
ChatBskyConvoDefs.isLogCreateMessage(log) &&
|
||||||
ChatBskyConvoDefs.isMessageView(log.message)
|
ChatBskyConvoDefs.isMessageView(log.message)
|
||||||
) {
|
) {
|
||||||
if (this.newMessages.has(log.message.id)) {
|
if (this.newMessages.has(log.message.id)) {
|
||||||
// Trust the log as the source of truth on ordering
|
// Trust the log as the source of truth on ordering
|
||||||
this.newMessages.delete(log.message.id)
|
this.newMessages.delete(log.message.id)
|
||||||
}
|
}
|
||||||
this.newMessages.set(log.message.id, log.message)
|
this.newMessages.set(log.message.id, log.message)
|
||||||
needsCommit = true
|
needsCommit = true
|
||||||
} else if (
|
} else if (
|
||||||
ChatBskyConvoDefs.isLogDeleteMessage(log) &&
|
ChatBskyConvoDefs.isLogDeleteMessage(log) &&
|
||||||
ChatBskyConvoDefs.isDeletedMessageView(log.message)
|
ChatBskyConvoDefs.isDeletedMessageView(log.message)
|
||||||
) {
|
) {
|
||||||
|
/*
|
||||||
|
* Update if we have this in state. If we don't, don't worry about it.
|
||||||
|
*/
|
||||||
|
if (this.pastMessages.has(log.message.id)) {
|
||||||
/*
|
/*
|
||||||
* Update if we have this in state. If we don't, don't worry about it.
|
* For now, we remove deleted messages from the thread, if we receive one.
|
||||||
|
*
|
||||||
|
* To support them, it'd look something like this:
|
||||||
|
* this.pastMessages.set(log.message.id, log.message)
|
||||||
*/
|
*/
|
||||||
if (this.pastMessages.has(log.message.id)) {
|
this.pastMessages.delete(log.message.id)
|
||||||
/*
|
this.newMessages.delete(log.message.id)
|
||||||
* For now, we remove deleted messages from the thread, if we receive one.
|
this.deletedMessages.delete(log.message.id)
|
||||||
*
|
needsCommit = true
|
||||||
* To support them, it'd look something like this:
|
|
||||||
* this.pastMessages.set(log.message.id, log.message)
|
|
||||||
*/
|
|
||||||
this.pastMessages.delete(log.message.id)
|
|
||||||
this.newMessages.delete(log.message.id)
|
|
||||||
this.deletedMessages.delete(log.message.id)
|
|
||||||
needsCommit = true
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
if (needsCommit) {
|
if (needsCommit) {
|
||||||
this.commit()
|
|
||||||
}
|
|
||||||
} catch (e: any) {
|
|
||||||
logger.error('Convo: failed to poll events')
|
|
||||||
this.pollingFailure = true
|
|
||||||
this.footerItems.set(ConvoItemError.PollFailed, {
|
|
||||||
type: 'error-recoverable',
|
|
||||||
key: ConvoItemError.PollFailed,
|
|
||||||
code: ConvoItemError.PollFailed,
|
|
||||||
retry: () => {
|
|
||||||
this.footerItems.delete(ConvoItemError.PollFailed)
|
|
||||||
this.pollingFailure = false
|
|
||||||
this.commit()
|
|
||||||
this.pollEvents()
|
|
||||||
},
|
|
||||||
})
|
|
||||||
this.commit()
|
this.commit()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
import React, {useContext, useState, useSyncExternalStore} from 'react'
|
import React, {useContext, useState, useSyncExternalStore} from 'react'
|
||||||
|
import {AppState} from 'react-native'
|
||||||
import {BskyAgent} from '@atproto-labs/api'
|
import {BskyAgent} from '@atproto-labs/api'
|
||||||
import {useFocusEffect} from '@react-navigation/native'
|
import {useFocusEffect} from '@react-navigation/native'
|
||||||
|
|
||||||
@@ -44,5 +45,21 @@ export function ChatProvider({
|
|||||||
}, [convo]),
|
}, [convo]),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
React.useEffect(() => {
|
||||||
|
const handleAppStateChange = (nextAppState: string) => {
|
||||||
|
if (nextAppState === 'active') {
|
||||||
|
convo.resume()
|
||||||
|
} else {
|
||||||
|
convo.background()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
const sub = AppState.addEventListener('change', handleAppStateChange)
|
||||||
|
|
||||||
|
return () => {
|
||||||
|
sub.remove()
|
||||||
|
}
|
||||||
|
}, [convo])
|
||||||
|
|
||||||
return <ChatContext.Provider value={service}>{children}</ChatContext.Provider>
|
return <ChatContext.Provider value={service}>{children}</ChatContext.Provider>
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user