|
13 | 13 |
|
14 | 14 | 'use strict'; |
15 | 15 |
|
| 16 | +const async = require(`async`); |
16 | 17 | const pubsub = require(`@google-cloud/pubsub`)(); |
17 | 18 | const uuid = require(`node-uuid`); |
18 | 19 | const path = require(`path`); |
@@ -58,57 +59,66 @@ describe(`pubsub:topics`, () => { |
58 | 59 | }); |
59 | 60 |
|
60 | 61 | it(`should publish a simple message`, (done) => { |
61 | | - pubsub.topic(topicName).subscribe(subscriptionName, (err, subscription) => { |
62 | | - assert.ifError(err); |
63 | | - run(`${cmd} publish ${topicName} "${message.data}"`, cwd); |
64 | | - setTimeout(() => { |
65 | | - subscription.pull((err, messages) => { |
66 | | - assert.ifError(err); |
67 | | - assert.equal(messages[0].data, message.data); |
68 | | - done(); |
69 | | - }); |
70 | | - }, 2000); |
71 | | - }); |
| 62 | + async.waterfall([ |
| 63 | + (cb) => { |
| 64 | + pubsub.topic(topicName).subscribe(subscriptionName, cb); |
| 65 | + }, |
| 66 | + (subscription, apiResponse, cb) => { |
| 67 | + run(`${cmd} publish ${topicName} "${message.data}"`, cwd); |
| 68 | + setTimeout(() => subscription.pull(cb), 2000); |
| 69 | + }, |
| 70 | + (messages, apiResponse, cb) => { |
| 71 | + assert.equal(messages[0].data, message.data); |
| 72 | + cb(); |
| 73 | + } |
| 74 | + ], done); |
72 | 75 | }); |
73 | 76 |
|
74 | 77 | it(`should publish a JSON message`, (done) => { |
75 | | - pubsub.topic(topicName).subscribe(subscriptionName, { reuseExisting: true }, (err, subscription) => { |
76 | | - assert.ifError(err); |
77 | | - run(`${cmd} publish ${topicName} '${JSON.stringify(message)}'`, cwd); |
78 | | - setTimeout(() => { |
79 | | - subscription.pull((err, messages) => { |
80 | | - assert.ifError(err); |
81 | | - assert.deepEqual(messages[0].data, message); |
82 | | - done(); |
83 | | - }); |
84 | | - }, 2000); |
85 | | - }); |
| 78 | + async.waterfall([ |
| 79 | + (cb) => { |
| 80 | + pubsub.topic(topicName).subscribe(subscriptionName, { reuseExisting: true }, cb); |
| 81 | + }, |
| 82 | + (subscription, apiResponse, cb) => { |
| 83 | + run(`${cmd} publish ${topicName} '${JSON.stringify(message)}'`, cwd); |
| 84 | + setTimeout(() => subscription.pull(cb), 2000); |
| 85 | + }, |
| 86 | + (messages, apiResponse, cb) => { |
| 87 | + assert.deepEqual(messages[0].data, message); |
| 88 | + cb(); |
| 89 | + } |
| 90 | + ], done); |
86 | 91 | }); |
87 | 92 |
|
88 | 93 | it(`should publish ordered messages`, (done) => { |
89 | 94 | const topics = require('../topics'); |
90 | | - pubsub.topic(topicName).subscribe(subscriptionName, { reuseExisting: true }, (err, subscription) => { |
91 | | - assert.ifError(err); |
92 | | - topics.publishOrderedMessage(topicName, message.data, () => { |
93 | | - setTimeout(() => { |
94 | | - subscription.pull((err, messages) => { |
95 | | - assert.ifError(err); |
96 | | - assert.equal(messages[0].data, message.data); |
97 | | - assert.equal(messages[0].attributes.orderId, '1'); |
98 | | - topics.publishOrderedMessage(topicName, message.data, () => { |
99 | | - setTimeout(() => { |
100 | | - subscription.pull((err, messages) => { |
101 | | - assert.ifError(err); |
102 | | - assert.equal(messages[0].data, message.data); |
103 | | - assert.equal(messages[0].attributes.orderId, '2'); |
104 | | - done(); |
105 | | - }); |
106 | | - }, 2000); |
107 | | - }); |
108 | | - }); |
109 | | - }, 2000); |
110 | | - }); |
111 | | - }); |
| 95 | + let subscription; |
| 96 | + |
| 97 | + async.waterfall([ |
| 98 | + (cb) => { |
| 99 | + pubsub.topic(topicName).subscribe(subscriptionName, { reuseExisting: true }, cb); |
| 100 | + }, |
| 101 | + (_subscription, apiResponse, cb) => { |
| 102 | + subscription = _subscription; |
| 103 | + topics.publishOrderedMessage(topicName, message.data, cb); |
| 104 | + }, |
| 105 | + (cb) => { |
| 106 | + setTimeout(() => subscription.pull(cb), 2000); |
| 107 | + }, |
| 108 | + (messages, apiResponse, cb) => { |
| 109 | + assert.equal(messages[0].data, message.data); |
| 110 | + assert.equal(messages[0].attributes.counterId, '1'); |
| 111 | + topics.publishOrderedMessage(topicName, message.data, cb); |
| 112 | + }, |
| 113 | + (cb) => { |
| 114 | + setTimeout(() => subscription.pull(cb), 2000); |
| 115 | + }, |
| 116 | + (messages, apiResponse, cb) => { |
| 117 | + assert.equal(messages[0].data, message.data); |
| 118 | + assert.equal(messages[0].attributes.counterId, '2'); |
| 119 | + topics.publishOrderedMessage(topicName, message.data, cb); |
| 120 | + } |
| 121 | + ], done); |
112 | 122 | }); |
113 | 123 |
|
114 | 124 | it(`should set the IAM policy for a topic`, (done) => { |
|
0 commit comments