mirror of
https://github.com/wassname/talk.git
synced 2026-09-09 11:38:08 +08:00
Cleanup of subscriptions
This commit is contained in:
+3
-32
@@ -1,16 +1,9 @@
|
||||
const schema = require('./schema');
|
||||
const Context = require('./context');
|
||||
const setupFunctions = require('./setupFunctions');
|
||||
|
||||
const {connectionOptions} = require('../services/redis');
|
||||
const {RedisPubSub} = require('graphql-redis-subscriptions');
|
||||
const {SubscriptionManager} = require('graphql-subscriptions');
|
||||
const {SubscriptionServer} = require('subscriptions-transport-ws');
|
||||
|
||||
const pubsub = new RedisPubSub(connectionOptions);
|
||||
const pubsub = require('./pubsub');
|
||||
const {createSubscriptionManager} = require('./subscriptions');
|
||||
|
||||
module.exports = {
|
||||
pubsub,
|
||||
createGraphOptions: (req) => ({
|
||||
|
||||
// Schema is created already, so just include it.
|
||||
@@ -20,27 +13,5 @@ module.exports = {
|
||||
// the lifespan of this request.
|
||||
context: new Context(req, pubsub)
|
||||
}),
|
||||
createSubscriptionManager: (server, path, sessionFactory) => new SubscriptionServer({
|
||||
subscriptionManager: new SubscriptionManager({
|
||||
schema,
|
||||
pubsub,
|
||||
setupFunctions,
|
||||
}),
|
||||
onSubscribe: (parsedMessage, baseParams, connection) => {
|
||||
|
||||
// Attach the context per request.
|
||||
baseParams.context = () => sessionFactory(connection.upgradeReq)
|
||||
.then((req) => new Context(req, pubsub))
|
||||
.catch((err) => {
|
||||
console.error(err);
|
||||
|
||||
return new Context({}, pubsub);
|
||||
});
|
||||
|
||||
return baseParams;
|
||||
}
|
||||
}, {
|
||||
server,
|
||||
path
|
||||
})
|
||||
createSubscriptionManager
|
||||
};
|
||||
|
||||
@@ -0,0 +1,5 @@
|
||||
const {RedisPubSub} = require('graphql-redis-subscriptions');
|
||||
|
||||
const {connectionOptions} = require('../services/redis');
|
||||
|
||||
module.exports = new RedisPubSub(connectionOptions);
|
||||
@@ -1,21 +0,0 @@
|
||||
const plugins = require('../services/plugins');
|
||||
const _ = require('lodash');
|
||||
|
||||
// Core setup functions
|
||||
let setupFunctions = {
|
||||
commentAdded: (options, args) => ({
|
||||
commentAdded: {
|
||||
filter: (comment) => comment.asset_id === args.asset_id
|
||||
},
|
||||
}),
|
||||
};
|
||||
|
||||
/**
|
||||
* Plugin support requires that we merge in existing setupFunctions with our new
|
||||
* plugin based ones. This allows plugins to extend existing setupFunctions as well
|
||||
* as provide new ones.
|
||||
*/
|
||||
module.exports = plugins.get('server', 'setupFunctions').reduce((acc, {setupFunctions}) => {
|
||||
|
||||
return _.merge(acc, setupFunctions);
|
||||
}, setupFunctions);
|
||||
@@ -0,0 +1,60 @@
|
||||
const {SubscriptionManager} = require('graphql-subscriptions');
|
||||
const {SubscriptionServer} = require('subscriptions-transport-ws');
|
||||
const _ = require('lodash');
|
||||
|
||||
const pubsub = require('./pubsub');
|
||||
const schema = require('./schema');
|
||||
const Context = require('./context');
|
||||
const plugins = require('../services/plugins');
|
||||
|
||||
const {deserializeUser} = require('../services/subscriptions');
|
||||
|
||||
// Core setup functions
|
||||
let setupFunctions = {
|
||||
commentAdded: (options, args) => ({
|
||||
commentAdded: {
|
||||
filter: (comment) => comment.asset_id === args.asset_id
|
||||
},
|
||||
}),
|
||||
};
|
||||
|
||||
/**
|
||||
* Plugin support requires that we merge in existing setupFunctions with our new
|
||||
* plugin based ones. This allows plugins to extend existing setupFunctions as well
|
||||
* as provide new ones.
|
||||
*/
|
||||
setupFunctions = plugins.get('server', 'setupFunctions').reduce((acc, {setupFunctions}) => {
|
||||
|
||||
return _.merge(acc, setupFunctions);
|
||||
}, setupFunctions);
|
||||
|
||||
/**
|
||||
* This creates a new subscription manager.
|
||||
*/
|
||||
const createSubscriptionManager = (server) => new SubscriptionServer({
|
||||
subscriptionManager: new SubscriptionManager({
|
||||
schema,
|
||||
pubsub,
|
||||
setupFunctions,
|
||||
}),
|
||||
onSubscribe: (parsedMessage, baseParams, connection) => {
|
||||
|
||||
// Attach the context per request.
|
||||
baseParams.context = () => deserializeUser(connection.upgradeReq)
|
||||
.then((req) => new Context(req, pubsub))
|
||||
.catch((err) => {
|
||||
console.error(err);
|
||||
|
||||
return new Context({}, pubsub);
|
||||
});
|
||||
|
||||
return baseParams;
|
||||
}
|
||||
}, {
|
||||
server,
|
||||
path: '/api/v1/live'
|
||||
});
|
||||
|
||||
module.exports = {
|
||||
createSubscriptionManager
|
||||
};
|
||||
Reference in New Issue
Block a user