-
Notifications
You must be signed in to change notification settings - Fork 161
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
when sending in a bench type of situation where buffers are large and…
… loops are tight, the libuv buffers for the socket can be exhausted. Changed the way that the outbound is serialized: - client will try to send 16K buffer at most - if the write returns that it buffered outside of kernel, client won't send until socket is drained
- Loading branch information
Showing
5 changed files
with
244 additions
and
30 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,64 @@ | ||
/* | ||
* Copyright 2020 The NATS Authors | ||
* Licensed under the Apache License, Version 2.0 (the "License"); | ||
* you may not use this file except in compliance with the License. | ||
* You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
|
||
/* jslint node: true */ | ||
'use strict' | ||
|
||
// change '../' to 'nats' when copying | ||
const nats = require('../') | ||
const argv = require('minimist')(process.argv.slice(2)) | ||
|
||
const url = argv.s || nats.DEFAULT_URI | ||
const creds = argv.creds | ||
const subject = argv._[0] | ||
const count = argv.c || 1 | ||
const len = argv.p || 0 | ||
let msg = argv._[1] || '' | ||
|
||
if (!subject) { | ||
console.log('Usage: node-pubsub [-s server] [--creds=filepath] <subject> [msg]') | ||
process.exit() | ||
} | ||
|
||
if (len) { | ||
for (let i = 0; i < len; i++) { | ||
msg += (i % 10) + '' | ||
} | ||
} | ||
|
||
// Connect to NATS server. | ||
const opts = nats.creds(creds) || {} | ||
opts.yieldTime = 1000 | ||
opts.name = 'bulk-pub' | ||
const nc = nats.connect(url, opts) | ||
.on('connect', () => { | ||
(async () => { | ||
let i = 0 | ||
for (i = 0; i < count; i++) { | ||
nc.publish(subject, msg) | ||
if (i % 1000 === 0) { | ||
console.info(`< ${i} messages`) | ||
} | ||
} | ||
nc.flush(() => { | ||
console.log(`Published ${i} messages`) | ||
}) | ||
})() | ||
}) | ||
|
||
nc.on('error', (e) => { | ||
console.log('Error [' + nc.currentServer + ']: ' + e) | ||
process.exit() | ||
}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,56 @@ | ||
/* | ||
* Copyright 2020 The NATS Authors | ||
* Licensed under the Apache License, Version 2.0 (the "License"); | ||
* you may not use this file except in compliance with the License. | ||
* You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
|
||
/* jslint node: true */ | ||
'use strict' | ||
|
||
// change '../' to 'nats' when copying | ||
const nats = require('../') | ||
const argv = require('minimist')(process.argv.slice(2)) | ||
|
||
const url = argv.s || nats.DEFAULT_URI | ||
const creds = argv.creds | ||
const subject = argv._[0] | ||
|
||
if (!subject) { | ||
console.log('Usage: bulk-sub [-s server] [--creds=filepath] <subject>') | ||
process.exit() | ||
} | ||
|
||
// Connect to NATS server. | ||
const opts = nats.creds(creds) || {} | ||
opts.yieldTime = 1000 | ||
const nc = nats.connect(url, opts) | ||
.on('connect', () => { | ||
console.log('connected'); | ||
(async () => { | ||
let i = 0 | ||
nc.subscribe(subject, (msg) => { | ||
i++ | ||
if (i % 1000 === 0) { | ||
console.info(`> ${i} messages`) | ||
} | ||
}) | ||
})() | ||
}) | ||
|
||
nc.on('close', () => { | ||
console.log('client closed') | ||
}) | ||
|
||
nc.on('error', (e) => { | ||
console.log('Error [' + nc.currentServer + ']: ' + e) | ||
process.exit() | ||
}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,69 @@ | ||
#!/usr/bin/env node | ||
|
||
/* | ||
* Copyright 2013-2020 The NATS Authors | ||
* Licensed under the Apache License, Version 2.0 (the "License"); | ||
* you may not use this file except in compliance with the License. | ||
* You may obtain a copy of the License at | ||
* | ||
* http://www.apache.org/licenses/LICENSE-2.0 | ||
* | ||
* Unless required by applicable law or agreed to in writing, software | ||
* distributed under the License is distributed on an "AS IS" BASIS, | ||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
* See the License for the specific language governing permissions and | ||
* limitations under the License. | ||
*/ | ||
|
||
/* jslint node: true */ | ||
'use strict' | ||
|
||
// change '../' to 'nats' when copying | ||
const nats = require('../') | ||
const argv = require('minimist')(process.argv.slice(2)) | ||
|
||
const url = argv.s || nats.DEFAULT_URI | ||
const creds = argv.creds | ||
const subject = argv._[0] | ||
const count = argv.c || 1 | ||
const msg = argv._[1] || '' | ||
|
||
if (!subject) { | ||
console.log('Usage: node-pub [-s server] [--creds=filepath] <subject> [msg]') | ||
process.exit() | ||
} | ||
|
||
// Connect to NATS server. | ||
const opts = nats.creds(creds) || {} | ||
opts.yieldTime = 1000 | ||
const nc = nats.connect(url, opts) | ||
.on('connect', () => { | ||
(async () => { | ||
let i = 0 | ||
nc.subscribe(subject, (msg) => { | ||
i++ | ||
if (i % 1000 === 0) { | ||
console.info(`> ${i} messages`) | ||
} | ||
}) | ||
})(); | ||
|
||
(async () => { | ||
let i = 0 | ||
for (i = 0; i < count; i++) { | ||
nc.publish(subject, msg) | ||
if (i % 1000 === 0) { | ||
console.info(`< ${i} messages`) | ||
} | ||
} | ||
nc.flush(() => { | ||
console.log(`Published ${i} messages`) | ||
process.exit() | ||
}) | ||
})() | ||
}) | ||
|
||
nc.on('error', (e) => { | ||
console.log('Error [' + nc.currentServer + ']: ' + e) | ||
process.exit() | ||
}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters