@@ -21,6 +21,7 @@ abstract class PubnubCoreAsync extends PubnubCore implements PubnubAsyncInterfac
2121
2222 protected TimedTaskManager timedTaskManager ;
2323 private volatile String _timetoken = "0" ;
24+ private volatile String _region = null ;
2425 private volatile String _saved_timetoken = "0" ;
2526
2627 protected static String PRESENCE_SUFFIX = "-pnpres" ;
@@ -33,7 +34,12 @@ abstract class PubnubCoreAsync extends PubnubCore implements PubnubAsyncInterfac
3334 private int HEARTBEAT = 320 ;
3435 private volatile int PRESENCE_HB_INTERVAL = 0 ;
3536
37+ private boolean V2 = false ;
3638
39+ public void setV2 (boolean v2 ) {
40+ this .V2 = v2 ;
41+ }
42+
3743 public void shutdown () {
3844 nonSubscribeManager .stop ();
3945 subscribeManager .stop ();
@@ -379,6 +385,103 @@ public void publish(String channel, Double message, Callback callback) {
379385 _publish (args , false );
380386 }
381387
388+ public void publish (String channel , JSONObject message , boolean storeInHistory , JSONObject metadata ,
389+ Callback callback ) {
390+ Hashtable args = new Hashtable ();
391+ args .put ("channel" , channel );
392+ args .put ("message" , message );
393+ args .put ("callback" , callback );
394+ args .put ("meta" , metadata );
395+ args .put ("storeInHistory" , (storeInHistory ) ? "" : "0" );
396+ _publish (args , false );
397+ }
398+
399+ public void publish (String channel , JSONArray message , boolean storeInHistory , JSONObject metadata ,
400+ Callback callback ) {
401+ Hashtable args = new Hashtable ();
402+ args .put ("channel" , channel );
403+ args .put ("message" , message );
404+ args .put ("callback" , callback );
405+ args .put ("meta" , metadata );
406+ args .put ("storeInHistory" , (storeInHistory ) ? "" : "0" );
407+ _publish (args , false );
408+ }
409+
410+ public void publish (String channel , String message , boolean storeInHistory , JSONObject metadata , Callback callback ) {
411+ Hashtable args = new Hashtable ();
412+ args .put ("channel" , channel );
413+ args .put ("message" , message );
414+ args .put ("callback" , callback );
415+ args .put ("meta" , metadata );
416+ args .put ("storeInHistory" , (storeInHistory ) ? "" : "0" );
417+ _publish (args , false );
418+ }
419+
420+ public void publish (String channel , Integer message , boolean storeInHistory , JSONObject metadata , Callback callback ) {
421+ Hashtable args = new Hashtable ();
422+ args .put ("channel" , channel );
423+ args .put ("message" , message );
424+ args .put ("callback" , callback );
425+ args .put ("meta" , metadata );
426+ args .put ("storeInHistory" , (storeInHistory ) ? "" : "0" );
427+ _publish (args , false );
428+ }
429+
430+ public void publish (String channel , Double message , boolean storeInHistory , JSONObject metadata , Callback callback ) {
431+ Hashtable args = new Hashtable ();
432+ args .put ("channel" , channel );
433+ args .put ("message" , message );
434+ args .put ("callback" , callback );
435+ args .put ("meta" , metadata );
436+ args .put ("storeInHistory" , (storeInHistory ) ? "" : "0" );
437+ _publish (args , false );
438+ }
439+
440+ public void publish (String channel , JSONObject message , JSONObject metadata , Callback callback ) {
441+ Hashtable args = new Hashtable ();
442+ args .put ("channel" , channel );
443+ args .put ("message" , message );
444+ args .put ("meta" , metadata );
445+ args .put ("callback" , callback );
446+ _publish (args , false );
447+ }
448+
449+ public void publish (String channel , JSONArray message , JSONObject metadata , Callback callback ) {
450+ Hashtable args = new Hashtable ();
451+ args .put ("channel" , channel );
452+ args .put ("message" , message );
453+ args .put ("meta" , metadata );
454+ args .put ("callback" , callback );
455+ _publish (args , false );
456+ }
457+
458+ public void publish (String channel , String message , JSONObject metadata , Callback callback ) {
459+ Hashtable args = new Hashtable ();
460+ args .put ("channel" , channel );
461+ args .put ("message" , message );
462+ args .put ("meta" , metadata );
463+ args .put ("callback" , callback );
464+ _publish (args , false );
465+ }
466+
467+ public void publish (String channel , Integer message , JSONObject metadata , Callback callback ) {
468+ Hashtable args = new Hashtable ();
469+ args .put ("channel" , channel );
470+ args .put ("message" , message );
471+ args .put ("meta" , metadata );
472+ args .put ("callback" , callback );
473+ _publish (args , false );
474+ }
475+
476+ public void publish (String channel , Double message , JSONObject metadata , Callback callback ) {
477+ Hashtable args = new Hashtable ();
478+ args .put ("channel" , channel );
479+ args .put ("message" , message );
480+ args .put ("meta" , metadata );
481+ args .put ("callback" , callback );
482+ _publish (args , false );
483+ }
484+
382485 protected void publish (Hashtable args , Callback callback ) {
383486 args .put ("callback" , callback );
384487 _publish (args , false );
@@ -1087,12 +1190,21 @@ private void _subscribe_base(boolean fresh, boolean dar, Worker worker) {
10871190 channelString = PubnubUtil .urlEncode (channelString );
10881191 }
10891192
1090- String [] urlComponents = { getPubnubUrl (), "subscribe" , this .SUBSCRIBE_KEY ,
1091- channelString , "0" + " /" + _timetoken };
1193+ String [] urlComponents = { getPubnubUrl (), (( this . V2 ) ? "v2/" : "" ) + "subscribe" , this .SUBSCRIBE_KEY ,
1194+ channelString , "0" + (( this . V2 ) ? "" : " /" + _timetoken ) };
10921195
10931196 Hashtable params = PubnubUtil .hashtableClone (this .params );
10941197 params .put ("uuid" , UUID );
10951198
1199+
1200+ if (this .V2 ) {
1201+ params .put ("tt" , _timetoken );
1202+ if (this ._region != null )
1203+ params .put ("tr" , this ._region );
1204+ } else {
1205+
1206+ }
1207+
10961208 if (groupsArray .length > 0 ) {
10971209 params .put ("channel-group" , groupString );
10981210 }
@@ -1106,8 +1218,79 @@ private void _subscribe_base(boolean fresh, boolean dar, Worker worker) {
11061218 log .verbose ("Subscribing with timetoken : " + _timetoken );
11071219
11081220
1221+ if (channelSubscriptions .getFilter () != null && channelSubscriptions .getFilter ().length () > 0 ) {
1222+ params .put ("filter-expr" , channelSubscriptions .getFilter ());
1223+ }
1224+
11091225 HttpRequest hreq = new HttpRequest (urlComponents , params , new ResponseHandler () {
11101226
1227+ void changeKey (JSONObject o , String ok , String nk ) throws JSONException {
1228+ if (!o .isNull (ok )) {
1229+ Object t = o .get (ok );
1230+ o .put (nk , t );
1231+ o .remove (ok );
1232+ }
1233+ }
1234+
1235+ JSONObject expandV2Keys (JSONObject m ) throws JSONException {
1236+ if (!m .isNull ("o" )) {
1237+ changeKey (m .getJSONObject ("o" ), "t" , "timetoken" );
1238+ changeKey (m .getJSONObject ("o" ), "r" , "region_code" );
1239+ }
1240+ if (!m .isNull ("p" )) {
1241+ changeKey (m .getJSONObject ("p" ), "t" , "timetoken" );
1242+ changeKey (m .getJSONObject ("p" ), "r" , "region_code" );
1243+ }
1244+ changeKey (m , "a" , "shard" );
1245+ changeKey (m , "b" , "subscription_match" );
1246+ changeKey (m , "c" , "channel" );
1247+ changeKey (m , "d" , "payload" );
1248+ changeKey (m , "ear" , "eat_after_reading" );
1249+ changeKey (m , "f" , "flags" );
1250+ changeKey (m , "i" , "issuing_client_id" );
1251+ changeKey (m , "k" , "subscribe_key" );
1252+ changeKey (m , "s" , "sequence_number" );
1253+ changeKey (m , "o" , "origination_timetoken" );
1254+ changeKey (m , "p" , "publish_timetoken" );
1255+ changeKey (m , "r" , "replication_map" );
1256+ changeKey (m , "u" , "user_metadata" );
1257+ changeKey (m , "w" , "waypoint_list" );
1258+ return m ;
1259+ }
1260+
1261+ void v2Handler (JSONObject jso , HttpRequest hreq ) throws JSONException {
1262+ JSONArray messages = jso .getJSONArray ("m" );
1263+ for (int i = 0 ; i < messages .length (); i ++) {
1264+ JSONObject messageObj = messages .getJSONObject (i );
1265+ String channel = messageObj .getString ("c" );
1266+ String sub_channel = (messageObj .isNull ("b" )) ? null : messageObj .getString ("b" );
1267+
1268+ String message = messageObj .getString ("d" );
1269+
1270+ SubscriptionItem chobj = null ;
1271+ if (channelSubscriptions != null && sub_channel != null )
1272+ chobj = channelSubscriptions .getItem (sub_channel );
1273+
1274+ if (chobj == null && channelGroupSubscriptions != null && sub_channel != null )
1275+ chobj = channelGroupSubscriptions .getItem (sub_channel );
1276+
1277+ if (chobj == null && channelSubscriptions != null )
1278+ chobj = channelSubscriptions .getItem (channel );
1279+
1280+ if (channel .indexOf ("-pnpres" ) > 0 ) {
1281+ chobj = channelSubscriptions .getItem (channel );
1282+ channel = PubnubUtil .splitString (channel , "-pnpres" )[0 ];
1283+
1284+ }
1285+
1286+ if (chobj != null ) {
1287+ Callback callback = chobj .callback ;
1288+ invokeSubscribeCallbackV2 (chobj .name , chobj .callback , message , expandV2Keys (messageObj ),
1289+ _timetoken , hreq );
1290+ }
1291+
1292+ }
1293+ }
11111294 void v1Handler (JSONArray jsa , HttpRequest hreq ) throws JSONException {
11121295
11131296 JSONArray messages = new JSONArray (jsa .get (0 ).toString ());
@@ -1162,19 +1345,29 @@ public void handleResponse(HttpRequest hreq, String response) {
11621345
11631346 String _in_response_timetoken = "" ;
11641347
1348+ boolean handleV2 = false ;
1349+
11651350 try {
11661351 jsa = new JSONArray (response );
11671352 _in_response_timetoken = jsa .get (1 ).toString ();
11681353
11691354 } catch (JSONException e ) {
1355+ try {
1356+ // handle V2 response
1357+ handleV2 = true ;
1358+ jso = new JSONObject (response );
11701359
1171- if (hreq .isSubzero ()) {
1172- log .verbose ("Response of subscribe 0 request. Need to do dAr process again" );
1173- _subscribe_base (false , hreq .isDar (), hreq .getWorker ());
1174- } else
1175- _subscribe_base (false );
1176- return ;
1360+ _in_response_timetoken = jso .getJSONObject ("t" ).getString ("t" );
1361+ _region = jso .getJSONObject ("t" ).getString ("r" );
11771362
1363+ } catch (JSONException e1 ) {
1364+ if (hreq .isSubzero ()) {
1365+ log .verbose ("Response of subscribe 0 request. Need to do dAr process again" );
1366+ _subscribe_base (false , hreq .isDar (), hreq .getWorker ());
1367+ } else
1368+ _subscribe_base (false );
1369+ return ;
1370+ }
11781371 }
11791372
11801373 /*
@@ -1203,7 +1396,10 @@ public void handleResponse(HttpRequest hreq, String response) {
12031396 }
12041397 try {
12051398
1206- v1Handler (jsa , hreq );
1399+ if (handleV2 )
1400+ v2Handler (jso , hreq );
1401+ else
1402+ v1Handler (jsa , hreq );
12071403
12081404 } catch (JSONException e ) {
12091405
@@ -1418,4 +1614,13 @@ public void unsetAuthKey() {
14181614 resubscribe ();
14191615 }
14201616
1617+
1618+ public String getFilter () {
1619+ return channelSubscriptions .getFilter ();
1620+ }
1621+
1622+ public void setFilter (String filter ) {
1623+ channelSubscriptions .setFilter (filter );
1624+ }
1625+
14211626}
0 commit comments