Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Submit feedback
Sign in
Toggle navigation
N
node-bigstream
Project
Project
Details
Activity
Releases
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
3
Merge Requests
3
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
bs
node-bigstream
Commits
ae3802f7
Commit
ae3802f7
authored
Apr 09, 2020
by
Kamron Aroonrua
💬
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
new rpc
parent
5cdbf970
Hide whitespace changes
Inline
Side-by-side
Showing
1 changed file
with
62 additions
and
0 deletions
+62
-0
rpccaller.js
lib/amqp/rpccaller.js
+62
-0
No files found.
lib/amqp/rpccaller.js
View file @
ae3802f7
var
amqp
=
require
(
'amqplib/callback_api'
);
var
amqp
=
require
(
'amqplib/callback_api'
);
const
REPLY_QUEUE
=
'amq.rabbitmq.reply-to'
;
function
RPCCaller
(
config
)
{
this
.
config
=
config
;
this
.
url
=
config
.
url
;
this
.
name
=
config
.
name
||
"rpc_queue"
;
this
.
conn
=
null
;
this
.
ch
=
null
;
var
self
=
this
;
this
.
opened
=
false
;
this
.
open
=
thunky
(
open
);
this
.
open
();
function
open
(
cb
)
{
amqp
.
connect
(
self
.
url
,
function
(
err
,
conn
)
{
if
(
err
){
return
cb
(
err
)}
conn
.
createChannel
(
function
(
err
,
ch
)
{
if
(
err
){
return
cb
(
err
)}
ch
.
responseEmitter
=
new
EventEmitter
();
ch
.
responseEmitter
.
setMaxListeners
(
0
);
ch
.
consume
(
REPLY_QUEUE
,
(
msg
)
=>
channel
.
responseEmitter
.
emit
(
msg
.
properties
.
correlationId
,
JSON
.
parse
(
msg
.
content
.
toString
())),
{
noAck
:
true
});
self
.
opened
=
true
;
self
.
conn
=
conn
;
self
.
ch
=
ch
;
cb
();
});
});
}
}
RPCCaller
.
prototype
.
call
=
function
(
req
,
cb
){
var
self
=
this
;
var
corr
=
generateUuid
();
this
.
open
(
function
(
err
){
if
(
err
){
console
.
log
(
err
);
}
self
.
ch
.
responseEmitter
.
once
(
corr
,
(
resp
)
=>
{
cb
(
null
,
resp
);
});
self
.
ch
.
sendToQueue
(
self
.
name
,
new
Buffer
(
JSON
.
stringify
(
req
)),
{
correlationId
:
corr
,
replyTo
:
REPLY_QUEUE
})
});
function
generateUuid
()
{
return
Math
.
random
().
toString
()
+
Math
.
random
().
toString
()
+
Math
.
random
().
toString
();
}
}
/*
function RPCCaller(config)
function RPCCaller(config)
{
{
this.config = config;
this.config = config;
...
@@ -35,5 +96,6 @@ RPCCaller.prototype.call = function(req,cb){
...
@@ -35,5 +96,6 @@ RPCCaller.prototype.call = function(req,cb){
Math.random().toString();
Math.random().toString();
}
}
}
}
*/
module
.
exports
=
RPCCaller
;
module
.
exports
=
RPCCaller
;
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment