Sitelet https://github.com/out0fmemory/java/commit/7630d9861313e4e6039563d6d6838181d385748d
Skip to content

Commit 7630d98

Browse files
author
Devendra
committed
adding psv2
1 parent cc8ce82 commit 7630d98

5 files changed

Lines changed: 3836 additions & 10 deletions

File tree

‎java/srcPubnubApi/srcCore/com/pubnub/api/PubnubCore.java‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -197,6 +197,11 @@ protected Object _publish(Hashtable args, boolean sync) {
197197

198198
if (storeInHistory != null && storeInHistory.length() > 0)
199199
parameters.put("store", storeInHistory);
200+
201+
JSONObject meta = (JSONObject) args.get("meta");
202+
if (meta != null && meta.length() > 0)
203+
parameters.put("meta", meta.toString());
204+
200205

201206
final Callback callback = getWrappedCallback(cb);
202207

‎java/srcPubnubApi/srcCore/com/pubnub/api/PubnubCoreAsync.java‎

Lines changed: 214 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -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
}

‎java/srcPubnubApi/srcCore/com/pubnub/api/Subscriptions.java‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,16 @@ class Subscriptions {
1515

1616
JSONObject state;
1717

18+
String filter;
19+
20+
public String getFilter() {
21+
return filter;
22+
}
23+
24+
public void setFilter(String filter) {
25+
this.filter = filter;
26+
}
27+
1828
void runConnectOnNewThread(final Callback callback, final String name, final JSONArray jsa) {
1929
Runnable r = new Runnable() {
2030
public void run() {

‎scala/scala-pubnub-tests/pom.xml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -166,7 +166,7 @@
166166
<reportsDirectory>${project.build.directory}/surefire-reports</reportsDirectory>
167167
<junitxml>.</junitxml>
168168
<filereports>WDF TestSuite.txt</filereports>
169-
<tagsToInclude>com.pubnub.api.tests.PublishTest</tagsToInclude>-->
169+
<!--<tagsToInclude>com.pubnub.api.tests.PublishTest</tagsToInclude>-->
170170
<systemProperties>
171171
<org.slf4j.simpleLogger.defaultLogLevel>debug</org.slf4j.simpleLogger.defaultLogLevel>
172172
</systemProperties>

0 commit comments

Comments
 (0)