Receive messages effectively Copy for LLM Open Markdown
JavaScript 13.0.3 Showing JavaScript examples. This guide shows you how to structure the subscribe side of a PubNub application. You'll learn how to:
Subscribe from a direct client connection instead of through a proxy.
Keep a checkpoint of the last message you processed.
Replay messages published while your client was disconnected, and skip the ones you already processed.
Copy received data into your own systems without forking every message from the client.
Every call on this page needs an SDK instance initialized with your subscribe key. If you don't have a keyset yet, start with Set up your account . This guide assumes you already know how to register a subscription. If you don't, start with Receive messages .
Examples use the JavaScript, Swift, Java, Kotlin, and Python SDKs, which document every parameter this guide uses. For any other language, refer to Available SDKs .
Subscribe from the client, not through a proxy
Open the subscription from the client that needs the messages, and let the SDK hold its own connection to PubNub's nearest point of presence. A proxy in the middle can conflict with the encryption and the long-lived TCP connection the SDK already maintains, and it inherits the proxy's own downtime as its own. For the reasoning behind a direct connection, refer to Architectural choices .
Recover messages published while you were disconnected
Live delivery to subscribers is at-most-once by default. On a stable connection, a subscriber receives each message at most once. After a reconnect, a replayed message can arrive again with the same timetoken. A subscriber can also miss messages if its buffer overflows or if it's disconnected when someone publishes a message.
To recover the gap, keep a checkpoint of the last message you processed and replay history from it after a reconnect. The recipe has four parts:
A durable checkpoint. Store the timetoken of the last message your handler finished, and write it only after the handler succeeds. Load it from storage when your app starts, so a restart recovers the same way a reconnect does.
A hold on live delivery. When the status listener reports a disconnect, buffer live messages instead of processing them.
Bounded, paged replay. When the client connects, page backward through history from now to the checkpoint, up to a page limit. If the gap needs more pages than the limit, reload state from your server instead of replaying.
One merge path. Process the replayed messages first, then the buffered live messages in timetoken order. Every message goes through the same check, which skips any message at or before the checkpoint, so a message that arrives on both paths runs once.
In each example, the checkpoint storage, the message handler, and the reload function are placeholders for your own code. The handler throws on failure, which leaves the checkpoint where it was and starts a replay from it.
JavaScript Swift Java Kotlin Python 1 const CHANNEL = 'channel_1' ; 2 const PAGE_SIZE = 100 ; 3 const MAX_PAGES = 10 ; 4 const RETRY_DELAY_MS = 5000 ; 5
6 7 let checkpoint = loadCheckpoint ( ) ; 8 let holdLive = checkpoint !== null ; 9 let recoveryRunning = false ; 10 const liveBuffer = [ ] ; 11
12 function processMessage ( timetoken , message ) { 13 if ( checkpoint !== null && BigInt ( timetoken ) <= BigInt ( checkpoint ) ) return ; 14 handleMessage ( message ) ; 15 checkpoint = timetoken ; 16 saveCheckpoint ( checkpoint ) ; 17 } 18
19 subscription . onMessage = ( event ) => { 20 const timetoken = String ( event . timetoken ) ; 21 if ( holdLive ) { 22 liveBuffer . push ( { timetoken , message : event . message } ) ; 23 return ; 24 } 25 try { 26 processMessage ( timetoken , event . message ) ; 27 } catch ( error ) { 28 console . error ( 'Handler failed, replaying from the checkpoint:' , error ) ; 29 recoverMissedMessages ( ) ; 30 } 31 } ; 32
33 pubnub . addListener ( { 34 status : ( event ) => { 35 switch ( event . category ) { 36 case 'PNDisconnectedCategory' : 37 case 'PNDisconnectedUnexpectedlyCategory' : 38 case 'PNConnectionErrorCategory' : 39 case 'PNNetworkDownCategory' : 40 holdLive = true ; 41 break ; 42 case 'PNConnectedCategory' : 43 case 'PNReconnectedCategory' : 44 recoverMissedMessages ( ) ; 45 break ; 46 } 47 } , 48 } ) ; 49
50 async function recoverMissedMessages ( ) { 51 if ( recoveryRunning ) return ; 52 recoveryRunning = true ; 53 holdLive = true ; 54 try { 55 if ( checkpoint !== null ) { 56 for ( const item of await fetchSinceCheckpoint ( ) ) { 57 processMessage ( String ( item . timetoken ) , item . message ) ; 58 } 59 } 60 liveBuffer . sort ( ( a , b ) => ( BigInt ( a . timetoken ) < BigInt ( b . timetoken ) ? - 1 : 1 ) ) ; 61 while ( liveBuffer . length > 0 ) { 62 processMessage ( liveBuffer [ 0 ] . timetoken , liveBuffer [ 0 ] . message ) ; 63 liveBuffer . shift ( ) ; 64 } 65 holdLive = false ; 66 recoveryRunning = false ; 67 } catch ( error ) { 68 recoveryRunning = false ; 69 if ( error . gapTooLarge ) { 70 reloadStateFromYourServer ( ) ; 71 return ; 72 } 73 console . error ( 'Recovery failed, retrying:' , error ) ; 74 setTimeout ( recoverMissedMessages , RETRY_DELAY_MS ) ; 75 } 76 } 77
78 79 async function fetchSinceCheckpoint ( ) { 80 const pages = [ ] ; 81 let start ; 82 for ( let page = 0 ; page < MAX_PAGES ; page ++ ) { 83 const response = await pubnub . fetchMessages ( { 84 channels : [ CHANNEL ] , 85 end : checkpoint , 86 start , 87 count : PAGE_SIZE , 88 } ) ; 89 const items = response . channels [ CHANNEL ] ?? [ ] ; 90 pages . unshift ( items ) ; 91 if ( items . length < PAGE_SIZE ) return pages . flat ( ) ; 92 start = String ( items [ 0 ] . timetoken ) ; 93 } 94 throw Object . assign ( new Error ( 'Gap exceeds MAX_PAGES' ) , { gapTooLarge : true } ) ; 95 } show all 95 lines 1 let channel = "channel_1" 2 let pageSize = 100 3 let maxPages = 10 4 let retryDelay : TimeInterval = 5 5
6 enum RecoveryError : Error { 7 case gapTooLarge 8 } 9
10 11 var checkpoint : Timetoken ? = loadCheckpoint ( ) 12 var holdLive = checkpoint != nil 13 var recoveryRunning = false 14 var liveBuffer : [ PubNubMessage ] = [ ] 15
16 func process ( _ message : PubNubMessage ) throws { 17 if let last = checkpoint , message . published <= last { return } 18 try handleMessage ( message ) 19 checkpoint = message . published 20 saveCheckpoint ( message . published ) 21 } 22
23 subscription . onMessage = { message in 24 if holdLive { 25 liveBuffer . append ( message ) 26 return 27 } 28 do { 29 try process ( message ) 30 } catch { 31 print ( "Handler failed, replaying from the checkpoint: \( error ) " ) 32 recoverMissedMessages ( ) 33 } 34 } 35
36 pubnub . onConnectionStateChange = { status in 37 switch status { 38 case . disconnected , . disconnectedUnexpectedly , . connectionError : 39 holdLive = true 40 case . connected : 41 recoverMissedMessages ( ) 42 default : 43 break 44 } 45 } 46
47 func recoverMissedMessages ( ) { 48 guard ! recoveryRunning else { return } 49 recoveryRunning = true 50 holdLive = true 51 guard let from = checkpoint else { 52 finishRecovery ( missed : [ ] ) 53 return 54 } 55 fetchPages ( end : from , start : nil , pages : [ ] ) { result in 56 switch result { 57 case let . success ( missed ) : 58 finishRecovery ( missed : missed ) 59 case . failure ( RecoveryError . gapTooLarge ) : 60 recoveryRunning = false 61 reloadStateFromYourServer ( ) 62 case let . failure ( error ) : 63 retryRecovery ( after : error ) 64 } 65 } 66 } 67
68 func finishRecovery ( missed : [ PubNubMessage ] ) { 69 do { 70 for message in missed { 71 try process ( message ) 72 } 73 liveBuffer . sort { $0 . published < $1 . published } 74 while let next = liveBuffer . first { 75 try process ( next ) 76 liveBuffer . removeFirst ( ) 77 } 78 holdLive = false 79 recoveryRunning = false 80 } catch { 81 retryRecovery ( after : error ) 82 } 83 } 84
85 func retryRecovery ( after error : Error ) { 86 print ( "Recovery failed, retrying: \( error ) " ) 87 recoveryRunning = false 88 DispatchQueue . main . asyncAfter ( deadline : . now ( ) + retryDelay ) { 89 recoverMissedMessages ( ) 90 } 91 } 92
93 94 func fetchPages ( 95 end : Timetoken , 96 start : Timetoken ? , 97 pages : [ [ PubNubMessage ] ] , 98 completion : @escaping ( Result < [ PubNubMessage ] , Error > ) -> Void 99 ) { 100 guard pages . count < maxPages else { 101 completion ( . failure ( RecoveryError . gapTooLarge ) ) 102 return 103 } 104 pubnub . fetchMessageHistory ( 105 for : [ channel ] , 106 page : PubNubBoundedPageBase ( start : start , end : end , limit : pageSize ) 107 ) { result in 108 switch result { 109 case let . success ( response ) : 110 let items = response . messagesByChannel [ channel ] ?? [ ] 111 let collected = [ items ] + pages 112 if items . count < pageSize { 113 completion ( . success ( collected . flatMap { $0 } ) ) 114 } else { 115 fetchPages ( end : end , start : items [ 0 ] . published , pages : collected , completion : completion ) 116 } 117 case let . failure ( error ) : 118 completion ( . failure ( error ) ) 119 } 120 } 121 } show all 121 lines 1 import com . google . gson . JsonElement ; 2 import com . pubnub . api . PubNubException ; 3 import com . pubnub . api . enums . PNStatusCategory ; 4 import com . pubnub . api . java . PubNub ; 5 import com . pubnub . api . java . v2 . subscriptions . Subscription ; 6 import com . pubnub . api . models . consumer . history . PNFetchMessageItem ; 7 import com . pubnub . api . models . consumer . history . PNFetchMessagesResult ; 8 import com . pubnub . api . models . consumer . pubsub . PNMessageResult ; 9
10 import java . util . ArrayList ; 11 import java . util . Collections ; 12 import java . util . Comparator ; 13 import java . util . LinkedList ; 14 import java . util . List ; 15 import java . util . concurrent . Executors ; 16 import java . util . concurrent . ScheduledExecutorService ; 17 import java . util . concurrent . TimeUnit ; 18
19 public abstract class MessageRecovery { 20 private static final String CHANNEL = "channel_1" ; 21 private static final int PAGE_SIZE = 100 ; 22 private static final int MAX_PAGES = 10 ; 23 private static final long RETRY_DELAY_MS = 5000 ; 24
25 private final PubNub pubnub ; 26 private final Object lock = new Object ( ) ; 27 private final List < PNMessageResult > liveBuffer = new ArrayList < > ( ) ; 28 private final ScheduledExecutorService executor = Executors . newSingleThreadScheduledExecutor ( ) ; 29 private Long checkpoint ; 30 private boolean holdLive ; 31 private boolean recoveryRunning = false ; 32
33 34 protected abstract Long loadCheckpoint ( ) ; 35 protected abstract void saveCheckpoint ( long timetoken ) ; 36 37 protected abstract void handleMessage ( JsonElement message ) ; 38 protected abstract void reloadStateFromYourServer ( ) ; 39
40 public MessageRecovery ( PubNub pubnub , Subscription subscription ) { 41 this . pubnub = pubnub ; 42 checkpoint = loadCheckpoint ( ) ; 43 holdLive = checkpoint != null ; 44 subscription . setOnMessage ( this :: onMessage ) ; 45 pubnub . addListener ( ( pn , status ) -> onStatus ( status . getCategory ( ) ) ) ; 46 } 47
48 private void processMessage ( long timetoken , JsonElement message ) { 49 if ( checkpoint != null && timetoken <= checkpoint ) return ; 50 handleMessage ( message ) ; 51 checkpoint = timetoken ; 52 saveCheckpoint ( timetoken ) ; 53 } 54
55 private void onMessage ( PNMessageResult event ) { 56 synchronized ( lock ) { 57 if ( holdLive ) { 58 liveBuffer . add ( event ) ; 59 return ; 60 } 61 try { 62 processMessage ( event . getTimetoken ( ) , event . getMessage ( ) ) ; 63 } catch ( RuntimeException error ) { 64 System . err . println ( "Handler failed, replaying from the checkpoint: " + error ) ; 65 startRecovery ( ) ; 66 } 67 } 68 } 69
70 private void onStatus ( PNStatusCategory category ) { 71 switch ( category ) { 72 case PNDisconnectedCategory : 73 case PNUnexpectedDisconnectCategory : 74 case PNConnectionError : 75 synchronized ( lock ) { 76 holdLive = true ; 77 } 78 break ; 79 case PNConnectedCategory : 80 startRecovery ( ) ; 81 break ; 82 default : 83 break ; 84 } 85 } 86
87 private void startRecovery ( ) { 88 synchronized ( lock ) { 89 if ( recoveryRunning ) return ; 90 recoveryRunning = true ; 91 holdLive = true ; 92 } 93 executor . execute ( this :: recoverMissedMessages ) ; 94 } 95
96 private void recoverMissedMessages ( ) { 97 try { 98 List < PNFetchMessageItem > missed = fetchSinceCheckpoint ( ) ; 99 synchronized ( lock ) { 100 for ( PNFetchMessageItem item : missed ) { 101 processMessage ( item . getTimetoken ( ) , item . getMessage ( ) ) ; 102 } 103 liveBuffer . sort ( Comparator . comparing ( PNMessageResult :: getTimetoken ) ) ; 104 while ( ! liveBuffer . isEmpty ( ) ) { 105 PNMessageResult event = liveBuffer . get ( 0 ) ; 106 processMessage ( event . getTimetoken ( ) , event . getMessage ( ) ) ; 107 liveBuffer . remove ( 0 ) ; 108 } 109 holdLive = false ; 110 recoveryRunning = false ; 111 } 112 } catch ( GapTooLargeException error ) { 113 synchronized ( lock ) { 114 recoveryRunning = false ; 115 } 116 reloadStateFromYourServer ( ) ; 117 } catch ( Exception error ) { 118 System . err . println ( "Recovery failed, retrying: " + error ) ; 119 synchronized ( lock ) { 120 recoveryRunning = false ; 121 } 122 executor . schedule ( this :: startRecovery , RETRY_DELAY_MS , TimeUnit . MILLISECONDS ) ; 123 } 124 } 125
126 127 private List < PNFetchMessageItem > fetchSinceCheckpoint ( ) throws PubNubException , GapTooLargeException { 128 Long end ; 129 synchronized ( lock ) { 130 end = checkpoint ; 131 } 132 if ( end == null ) return Collections . emptyList ( ) ; 133 LinkedList < List < PNFetchMessageItem > > pages = new LinkedList < > ( ) ; 134 Long start = null ; 135 for ( int page = 0 ; page < MAX_PAGES ; page ++ ) { 136 PNFetchMessagesResult response = pubnub . fetchMessages ( ) 137 . channels ( Collections . singletonList ( CHANNEL ) ) 138 . end ( end ) 139 . start ( start ) 140 . maximumPerChannel ( PAGE_SIZE ) 141 . sync ( ) ; 142 List < PNFetchMessageItem > items = response . getChannels ( ) . getOrDefault ( CHANNEL , Collections . emptyList ( ) ) ; 143 pages . addFirst ( items ) ; 144 if ( items . size ( ) < PAGE_SIZE ) { 145 List < PNFetchMessageItem > all = new ArrayList < > ( ) ; 146 pages . forEach ( all :: addAll ) ; 147 return all ; 148 } 149 start = items . get ( 0 ) . getTimetoken ( ) ; 150 } 151 throw new GapTooLargeException ( ) ; 152 } 153
154 private static class GapTooLargeException extends Exception { 155 } 156 } show all 156 lines 1 import com . google . gson . JsonElement 2 import com . pubnub . api . PubNub 3 import com . pubnub . api . enums . PNStatusCategory 4 import com . pubnub . api . models . consumer . PNBoundedPage 5 import com . pubnub . api . models . consumer . PNStatus 6 import com . pubnub . api . models . consumer . history . PNFetchMessageItem 7 import com . pubnub . api . models . consumer . pubsub . PNMessageResult 8 import com . pubnub . api . v2 . callbacks . StatusListener 9 import com . pubnub . api . v2 . subscriptions . Subscription 10 import java . util . concurrent . Executors 11 import java . util . concurrent . TimeUnit 12
13 abstract class MessageRecovery ( private val pubnub : PubNub , subscription : Subscription ) { 14 private val lock = Any ( ) 15 private val liveBuffer = mutableListOf < PNMessageResult > ( ) 16 private val executor = Executors . newSingleThreadScheduledExecutor ( ) 17 private var checkpoint : Long ? = null 18 private var holdLive = false 19 private var recoveryRunning = false 20
21 22 protected abstract fun loadCheckpoint ( ) : Long ? 23 protected abstract fun saveCheckpoint ( timetoken : Long ) 24 25 protected abstract fun handleMessage ( message : JsonElement ) 26 protected abstract fun reloadStateFromYourServer ( ) 27
28 init { 29 checkpoint = loadCheckpoint ( ) 30 holdLive = checkpoint != null 31 subscription . onMessage = { event -> onMessage ( event ) } 32 pubnub . addListener ( object : StatusListener { 33 override fun status ( pubnub : PubNub , status : PNStatus ) = onStatus ( status . category ) 34 } ) 35 } 36
37 private fun processMessage ( timetoken : Long , message : JsonElement ) { 38 checkpoint ? . let { if ( timetoken <= it ) return } 39 handleMessage ( message ) 40 checkpoint = timetoken 41 saveCheckpoint ( timetoken ) 42 } 43
44 private fun onMessage ( event : PNMessageResult ) = synchronized ( lock ) { 45 if ( holdLive ) { 46 liveBuffer . add ( event ) 47 return 48 } 49 try { 50 processMessage ( event . timetoken !! , event . message ) 51 } catch ( error : Exception ) { 52 println ( "Handler failed, replaying from the checkpoint: $ error " ) 53 startRecovery ( ) 54 } 55 } 56
57 private fun onStatus ( category : PNStatusCategory ) { 58 when ( category ) { 59 PNStatusCategory . PNDisconnectedCategory , 60 PNStatusCategory . PNUnexpectedDisconnectCategory , 61 PNStatusCategory . PNConnectionError -> synchronized ( lock ) { holdLive = true } 62 PNStatusCategory . PNConnectedCategory -> startRecovery ( ) 63 else -> Unit 64 } 65 } 66
67 private fun startRecovery ( ) { 68 synchronized ( lock ) { 69 if ( recoveryRunning ) return 70 recoveryRunning = true 71 holdLive = true 72 } 73 executor . execute { recoverMissedMessages ( ) } 74 } 75
76 private fun recoverMissedMessages ( ) { 77 try { 78 val missed = fetchSinceCheckpoint ( ) 79 synchronized ( lock ) { 80 missed . forEach { processMessage ( it . timetoken !! , it . message ) } 81 liveBuffer . sortBy { it . timetoken } 82 while ( liveBuffer . isNotEmpty ( ) ) { 83 val event = liveBuffer . first ( ) 84 processMessage ( event . timetoken !! , event . message ) 85 liveBuffer . removeAt ( 0 ) 86 } 87 holdLive = false 88 recoveryRunning = false 89 } 90 } catch ( error : GapTooLargeException ) { 91 synchronized ( lock ) { recoveryRunning = false } 92 reloadStateFromYourServer ( ) 93 } catch ( error : Exception ) { 94 println ( "Recovery failed, retrying: $ error " ) 95 synchronized ( lock ) { recoveryRunning = false } 96 executor . schedule ( { startRecovery ( ) } , RETRY_DELAY_MS , TimeUnit . MILLISECONDS ) 97 } 98 } 99
100 101 private fun fetchSinceCheckpoint ( ) : List < PNFetchMessageItem > { 102 val end = synchronized ( lock ) { checkpoint } ?: return emptyList ( ) 103 val pages = ArrayDeque < List < PNFetchMessageItem > > ( ) 104 var start : Long ? = null 105 repeat ( MAX_PAGES ) { 106 val response = pubnub . fetchMessages ( 107 channels = listOf ( CHANNEL ) , 108 page = PNBoundedPage ( start = start , end = end , limit = PAGE_SIZE ) , 109 ) . sync ( ) 110 val items = response . channels [ CHANNEL ] . orEmpty ( ) 111 pages . addFirst ( items ) 112 if ( items . size < PAGE_SIZE ) return pages . flatten ( ) 113 start = items . first ( ) . timetoken 114 } 115 throw GapTooLargeException ( ) 116 } 117
118 private class GapTooLargeException : Exception ( ) 119
120 companion object { 121 private const val CHANNEL = "channel_1" 122 private const val PAGE_SIZE = 100 123 private const val MAX_PAGES = 10 124 private const val RETRY_DELAY_MS = 5000L 125 } 126 } show all 126 lines 1 import threading 2
3 from pubnub . callbacks import SubscribeCallback 4 from pubnub . enums import PNStatusCategory 5
6 CHANNEL = 'channel_1' 7 PAGE_SIZE = 100 8 MAX_PAGES = 10 9 RETRY_DELAY_S = 5 10
11 12 checkpoint = load_checkpoint ( ) 13 hold_live = checkpoint is not None 14 recovery_running = False 15 live_buffer = [ ] 16 lock = threading . RLock ( ) 17
18
19 class GapTooLarge ( Exception ) : 20 pass 21
22
23 def process_message ( timetoken , message ) : 24 global checkpoint 25 if checkpoint is not None and timetoken <= checkpoint : 26 return 27 handle_message ( message ) 28 checkpoint = timetoken 29 save_checkpoint ( checkpoint ) 30
31
32 class RecoveryListener ( SubscribeCallback ) : 33 def message ( self , pubnub , event ) : 34 with lock : 35 if hold_live : 36 live_buffer . append ( ( int ( event . timetoken ) , event . message ) ) 37 return 38 try : 39 process_message ( int ( event . timetoken ) , event . message ) 40 except Exception as error : 41 print ( 'Handler failed, replaying from the checkpoint:' , error ) 42 start_recovery ( ) 43
44 def status ( self , pubnub , status ) : 45 global hold_live 46 if status . category in ( 47 PNStatusCategory . PNDisconnectedCategory , 48 PNStatusCategory . PNUnexpectedDisconnectCategory , 49 PNStatusCategory . PNConnectionErrorCategory , 50 ) : 51 with lock : 52 hold_live = True 53 elif status . category in ( 54 PNStatusCategory . PNConnectedCategory , 55 PNStatusCategory . PNReconnectedCategory , 56 ) : 57 start_recovery ( ) 58
59
60 def start_recovery ( ) : 61 global hold_live , recovery_running 62 with lock : 63 if recovery_running : 64 return 65 recovery_running = True 66 hold_live = True 67 threading . Thread ( target = recover_missed_messages , daemon = True ) . start ( ) 68
69
70 def recover_missed_messages ( ) : 71 global hold_live , recovery_running 72 try : 73 missed = fetch_since_checkpoint ( ) if checkpoint is not None else [ ] 74 with lock : 75 for timetoken , message in missed : 76 process_message ( timetoken , message ) 77 live_buffer . sort ( key = lambda item : item [ 0 ] ) 78 while live_buffer : 79 process_message ( * live_buffer [ 0 ] ) 80 live_buffer . pop ( 0 ) 81 hold_live = False 82 recovery_running = False 83 except GapTooLarge : 84 with lock : 85 recovery_running = False 86 reload_state_from_your_server ( ) 87 except Exception as error : 88 print ( 'Recovery failed, retrying:' , error ) 89 with lock : 90 recovery_running = False 91 threading . Timer ( RETRY_DELAY_S , start_recovery ) . start ( ) 92
93
94 95 def fetch_since_checkpoint ( ) : 96 pages = [ ] 97 start = None 98 for _ in range ( MAX_PAGES ) : 99 request = pubnub . fetch_messages ( ) . channels ( [ CHANNEL ] ) . end ( checkpoint ) . count ( PAGE_SIZE ) 100 if start is not None : 101 request = request . start ( start ) 102 items = request . sync ( ) . result . channels . get ( CHANNEL , [ ] ) 103 pages . insert ( 0 , [ ( int ( item . timetoken ) , item . message ) for item in items ] ) 104 if len ( items ) < PAGE_SIZE : 105 return [ item for page in pages for item in page ] 106 start = int ( items [ 0 ] . timetoken ) 107 raise GapTooLarge ( ) 108
109
110 pubnub . add_listener ( RecoveryListener ( ) ) show all 110 lines
The examples recover one channel. For several channels, keep one checkpoint and one live buffer per channel. The Swift example relies on the SDK's default of running callbacks on the main queue. If you set a different callback queue, serialize access to the recovery state, as the Java, Kotlin, and Python examples do with a lock.
The SDK retries a dropped connection on its own, and a short outage may not produce a disconnect status at all. When the client reconnects, the subscriber buffer delivers what it holds:
The subscriber message buffer queues messages for a reconnecting client. It holds 100 messages for up to 16 minutes by default, and discards the oldest first (FIFO) when a burst exceeds that size. Larger buffers, for example 300 or 500 messages, can be provisioned per keyset by PubNub Support.
The examples start recovery on a status event or a handler failure. If you don't want to rely on a status event for every gap, also call the recovery function on a schedule or when your app returns to the foreground. It's safe to call at any time, because the checkpoint check skips messages you already processed.
The checkpoint moves only after the handler returns. If your handler fails partway through, for example after it wrote to a database, it runs again on the same message, so make its side effects safe to repeat. For publish-side options, refer to Exactly-once processing .
For how to register the status listener in your SDK, refer to Monitor and respond to connection status changes . For the fetch call's parameters, page limits, and how far back it can reach, refer to Retrieve message history .
Copy received data to your own systems without forking every message
Don't have the client forward a second copy of every message to your server on top of handling it. That doubles the mobile data and battery cost of every message, on top of whatever the client already does with it.
Use an After Publish Function to forward a copy from PubNub's own network instead. Or read the data back later through Message Persistence , rather than capturing a copy of each message as it arrives. Refer to Copy published data to your own systems without publishing twice for both options. The choice is the same regardless of which side of the connection triggers it.