Merge branch 'next' into permalink

This commit is contained in:
Chi Vinh Le
2018-08-03 16:05:43 +02:00
127 changed files with 3521 additions and 1363 deletions
+5 -6
View File
@@ -3,12 +3,13 @@ import http from "http";
import { Redis } from "ioredis";
import { Db } from "mongodb";
import { Config } from "talk-common/config";
import { notFoundMiddleware } from "talk-server/app/middleware/notFound";
import { createPassport } from "talk-server/app/middleware/passport";
import { JWTSigningConfig } from "talk-server/app/middleware/passport/jwt";
import { Config } from "talk-server/config";
import { handleSubscriptions } from "talk-server/graph/common/subscriptions/middleware";
import { Schemas } from "talk-server/graph/schemas";
import TenantCache from "talk-server/services/tenant/cache";
import { accessLogger, errorLogger } from "./middleware/logging";
import serveStatic from "./middleware/serveStatic";
@@ -21,6 +22,7 @@ export interface AppOptions {
redis: Redis;
schemas: Schemas;
signingConfig: JWTSigningConfig;
tenantCache: TenantCache;
}
/**
@@ -34,10 +36,7 @@ export async function createApp(options: AppOptions): Promise<Express> {
parent.use(accessLogger);
// Create some services for the router.
const passport = createPassport({
db: options.mongo,
signingConfig: options.signingConfig,
});
const passport = createPassport(options);
// Mount the router.
parent.use(
@@ -76,7 +75,7 @@ export const listenAndServe = (
* handle websocket traffic by upgrading their http connections to websocket.
*
* @param schemas schemas for every schema this application handles
* @param server the http.Server to attach the websocket upgraders to
* @param server the http.Server to attach the websocket upgrader to
*/
export async function attachSubscriptionHandlers(
schemas: Schemas,
@@ -1,6 +1,6 @@
// Jest Snapshot v1, https://goo.gl/fbAQLP
exports[`createJWTSigningConfig parses a RSA certiciate 1`] = `
exports[`createJWTSigningConfig parses a RSA certificate 1`] = `
"-----BEGIN RSA PRIVATE KEY-----
MIIEpQIBAAKCAQEAyxR2DVlvkQRquggUQTpHN+PxDs2iOiItGgn6u4+faUCdgGEV
EnmG69//3lAZHnEQN9rkZS3/20zc41mTJnO7dslJbB316vWUSIwYcVY/VC9DTbk+
@@ -12,6 +12,7 @@ import { createLocalStrategy } from "talk-server/app/middleware/passport/local";
import { createOIDCStrategy } from "talk-server/app/middleware/passport/oidc";
import { createSSOStrategy } from "talk-server/app/middleware/passport/sso";
import { User } from "talk-server/models/user";
import TenantCache from "talk-server/services/tenant/cache";
import { Request } from "talk-server/types/express";
export type VerifyCallback = (
@@ -21,28 +22,28 @@ export type VerifyCallback = (
) => void;
export interface PassportOptions {
db: Db;
mongo: Db;
signingConfig: JWTSigningConfig;
tenantCache: TenantCache;
}
export function createPassport({
db,
signingConfig,
}: PassportOptions): passport.Authenticator {
export function createPassport(
options: PassportOptions
): passport.Authenticator {
// Create the authenticator.
const auth = new Authenticator();
// Use the OIDC Strategy.
auth.use(createOIDCStrategy({ db }));
auth.use(createOIDCStrategy(options));
// Use the LocalStrategy.
auth.use(createLocalStrategy({ db }));
auth.use(createLocalStrategy(options));
// Use the SSOStrategy.
auth.use(createSSOStrategy({ db }));
auth.use(createSSOStrategy(options));
// Use the JWTStrategy.
auth.use(createJWTStrategy({ db, signingConfig }));
auth.use(createJWTStrategy(options));
return auth;
}
@@ -1,11 +1,11 @@
import sinon from "sinon";
import { Config } from "talk-common/config";
import {
createJWTSigningConfig,
extractJWTFromRequest,
parseAuthHeader,
} from "talk-server/app/middleware/passport/jwt";
import { Config } from "talk-server/config";
import { Request } from "talk-server/types/express";
describe("parseAuthHeader", () => {
@@ -67,7 +67,7 @@ describe("extractJWTFromRequest", () => {
});
describe("createJWTSigningConfig", () => {
it("parses a RSA certiciate", () => {
it("parses a RSA certificate", () => {
const input = `-----BEGIN RSA PRIVATE KEY-----\\nMIIEpQIBAAKCAQEAyxR2DVlvkQRquggUQTpHN+PxDs2iOiItGgn6u4+faUCdgGEV\\nEnmG69//3lAZHnEQN9rkZS3/20zc41mTJnO7dslJbB316vWUSIwYcVY/VC9DTbk+\\nMHWZd94p5hOB8PoY2vEGA53KiyWLqQC5FWE3u7cz7eYTr9/eRPDTc15IzohLXd5U\\nC9EbO5ebho2CvWrBfrLozM5Kidp8r3Jp+A0o3kfJ/kRDDn/BmG6pM0TohWZFYMs2\\nnQaGg+of9tcafgAs7hZAgBrrcc/jke6+MKxpC8algik79nMk7s7prxF1Z9EbAeQV\\n1ssL2VgsjvGAHIV+Arckl6QJbVDvQXNAM0PqbQIDAQABAoIBAQCoG6D5vf5P8nMS\\n2ltB/6cyyfsjgO/45Y+mTXqERwj0DOwUeMkDyRv6KCxb8LxKade+FPIaG7D/7amw\\nfdcE7qrRUyD3YfnPbUk5oNcfAwFbg+BX969WWBMZmgvfDGj1fWKT4w9ScQ1YkFUD\\nKrkLzLVhK+/N0Dad0VjiguTXTMZCSDFOY9fO8HRF6EA3aewEPeEY62J6rSjGXvWB\\nGdW+FNvf/uRr36xGHNqiOP837pdVUppjgDyVsORnMfFtYMyWyxS2XD5r8gRwcRg7\\n0nz6bLM53DjKweO+Yl+pIVPFAyXL0pwzQDlnjShsCzyzjA9lJftkQwbcMWopeegJ\\nkPLmiq4VAoGBAOqDmySNx8vmWWMOaXKFuH6Gqu/Nd7gBHxZ73wvsEmvV52xwa0oi\\n55h+v6P1YEaNZQWXDFsvILoOUHr2kwZY+Du/MC7tgqpj+Fu3h7UHslulJRE3A+sN\\noLbHjZuwm3wwsatpHdyEYOGg0HIGWXi+9pDT/1gy8g3L2Gf0X6rfkBBXAoGBAN2v\\nlbii0+HvZ2y0D0P6NfUJ6cQDrSyuTe7UW6OVYjBjrVAk8+bhnQ4eKd9edCnUDqu6\\n9C8ZSrqR6VBeItbt8y+5ZCRcrigxd2VdH8rL9g6idD9RPnSbHx7Al8DxSUv25xMK\\n8Z/ZOAvuCmwDfdleycNDoTawKqLtWBzUEntLs5DbAoGAPlTKiJWylAxel8h92HWY\\nSvDqQCChgGOz6prz9sxBPS42e4kJy0OpwMt3jlGqzDXKswipvRayoSEq3PPqshY1\\nrFOtr9trDnTRzzbhuAkaq+ciCghQX0pY/BvgFJCFUyXyIzgmOrVotq+yl4v+fexr\\nxqTCSqQH2AjlNQQr5VPUi7MCgYEAsNbbMXE6YlXug+lS8CANoM3qm4FvSGA3LNhb\\nza9hp0YsP+1qXvgEp/lp35RiR+ewWE+HcHbVhOTWYFTnp9ojDyPtfZAtIUTsgIB7\\n1vNC8kOnRccSckQ32/k4VSJlHOL1S9yECMZnjiSyTZ2va5HQkyJE3PJE4LlCe6S0\\npYQq1tcCgYEAoJDeSeAPqi5NIu+MWNUWzw4vo5raKyHrJi+cTvKyM/2zJFHvBc5f\\nRaxkcIAOmIDoVdFgy6APY/0DnDnpqT1kMagUaxZjG9PLFIDds5DRaL99m+S7l8mt\\nySX/MbmhQHYWpVf2nL6pmfPuP4Ih6tbKIUUGA3wZXYYZ5r+pZFG1IrA=\\n-----END RSA PRIVATE KEY-----`;
const config = {
get: sinon.stub(),
+10 -11
View File
@@ -3,7 +3,7 @@ import uuid from "uuid";
import { Db } from "mongodb";
import { Strategy } from "passport-strategy";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
import { retrieveUser, User } from "talk-server/models/user";
import { Request } from "talk-server/types/express";
@@ -67,7 +67,7 @@ export function createAsymmetricSigningConfig(
secret: string
): JWTSigningConfig {
return {
// Secrets have their newlines encoded with newline litterals.
// Secrets have their newlines encoded with newline literals.
secret: Buffer.from(secret.replace(/\\n/g, "\n")),
algorithm,
};
@@ -138,21 +138,20 @@ export interface JWTToken {
export interface JWTStrategyOptions {
signingConfig: JWTSigningConfig;
db: Db;
mongo: Db;
}
export class JWTStrategy extends Strategy {
public name = "jwt";
private signingConfig: JWTSigningConfig;
private db: Db;
private mongo: Db;
public name: string;
constructor({ signingConfig, db }: JWTStrategyOptions) {
constructor({ signingConfig, mongo }: JWTStrategyOptions) {
super();
this.name = "jwt";
this.signingConfig = signingConfig;
this.db = db;
this.mongo = mongo;
}
public authenticate(req: Request) {
@@ -160,7 +159,7 @@ export class JWTStrategy extends Strategy {
const token = extractJWTFromRequest(req);
if (!token) {
// There was no token on the request, so there was no user, so let's mark
// that the strategy was succesfull.
// that the strategy was successful.
return this.success(null, null);
}
@@ -187,7 +186,7 @@ export class JWTStrategy extends Strategy {
try {
// Find the user.
const user = await retrieveUser(this.db, tenant.id, sub);
const user = await retrieveUser(this.mongo, tenant.id, sub);
// Return them! The user may be null, but that's ok here.
this.success(user, null);
@@ -8,7 +8,7 @@ import {
} from "talk-server/models/user";
import { Request } from "talk-server/types/express";
const verifyFactory = (db: Db) => async (
const verifyFactory = (mongo: Db) => async (
req: Request,
email: string,
password: string,
@@ -21,7 +21,7 @@ const verifyFactory = (db: Db) => async (
const tenant = req.tenant!;
// Get the user from the database.
const user = await retrieveUserWithProfile(db, tenant.id, {
const user = await retrieveUserWithProfile(mongo, tenant.id, {
id: email,
type: "local",
});
@@ -44,10 +44,10 @@ const verifyFactory = (db: Db) => async (
};
export interface LocalStrategyOptions {
db: Db;
mongo: Db;
}
export function createLocalStrategy({ db }: LocalStrategyOptions) {
export function createLocalStrategy({ mongo }: LocalStrategyOptions) {
return new LocalStrategy(
{
usernameField: "email",
@@ -55,6 +55,6 @@ export function createLocalStrategy({ db }: LocalStrategyOptions) {
session: false,
passReqToCallback: true,
},
verifyFactory(db)
verifyFactory(mongo)
);
}
+22 -16
View File
@@ -8,8 +8,10 @@ import { Strategy } from "passport-strategy";
import { validate } from "talk-server/app/request/body";
import { reconstructURL } from "talk-server/app/url";
import { GQLUSER_ROLE } from "talk-server/graph/tenant/schema/__generated__/types";
import { OIDCAuthIntegration, Tenant } from "talk-server/models/tenant";
import { OIDCAuthIntegration } from "talk-server/models/settings";
import { Tenant } from "talk-server/models/tenant";
import { OIDCProfile, retrieveUserWithProfile } from "talk-server/models/user";
import TenantCache from "talk-server/services/tenant/cache";
import { upsert } from "talk-server/services/users";
import { Request } from "talk-server/types/express";
@@ -39,10 +41,6 @@ export interface StrategyItem {
jwksClient?: JwksClient;
}
export interface OIDCStrategyOptions {
db: Db;
}
export function isOIDCToken(token: OIDCIDToken | object): token is OIDCIDToken {
if (
(token as OIDCIDToken).iss &&
@@ -176,20 +174,28 @@ export async function findOrCreateOIDCUser(
*/
const OIDC_SCOPE = "openid email profile";
// FIXME: attach strategy to cache updates of the tenants
export interface OIDCStrategyOptions {
mongo: Db;
tenantCache: TenantCache;
}
export default class OIDCStrategy extends Strategy {
public name: string;
public name = "oidc";
private db: Db;
private cache: Map<string, StrategyItem>;
private mongo: Db;
private cache = new Map<string, StrategyItem>();
constructor({ db }: OIDCStrategyOptions) {
constructor({ mongo, tenantCache }: OIDCStrategyOptions) {
super();
this.name = "oidc";
this.cache = new Map();
this.db = db;
this.mongo = mongo;
// Subscribe to updates with Tenants.
tenantCache.subscribe(tenant => {
// Delete the tenant cache item when the tenant changes. The refreshed
// Tenant will come in with the request.
this.cache.delete(tenant.id);
});
}
private lookupJWKSClient(
@@ -277,7 +283,7 @@ export default class OIDCStrategy extends Strategy {
try {
const user = await findOrCreateOIDCUser(
this.db,
this.mongo,
tenant,
decoded as OIDCIDToken
);
@@ -370,6 +376,6 @@ export default class OIDCStrategy extends Strategy {
}
}
export function createOIDCStrategy({ db }: OIDCStrategyOptions) {
return new OIDCStrategy({ db });
export function createOIDCStrategy(options: OIDCStrategyOptions) {
return new OIDCStrategy(options);
}
@@ -17,7 +17,7 @@ import { upsert } from "talk-server/services/users";
import { Request } from "talk-server/types/express";
export interface SSOStrategyOptions {
db: Db;
mongo: Db;
}
export interface SSOUserProfile {
@@ -114,15 +114,14 @@ export function isSSOToken(token: SSOToken | object): token is SSOToken {
}
export default class SSOStrategy extends Strategy {
public name: string;
public name = "sso";
private db: Db;
private mongo: Db;
constructor({ db }: SSOStrategyOptions) {
constructor({ mongo }: SSOStrategyOptions) {
super();
this.name = "sso";
this.db = db;
this.mongo = mongo;
}
/**
@@ -162,7 +161,7 @@ export default class SSOStrategy extends Strategy {
if (isOIDCToken(token)) {
// The token provided for SSO contains an issuer claim. We're assuming
// that this request is associated with an OpenID Connect provider.
return findOrCreateOIDCUser(this.db, tenant, token);
return findOrCreateOIDCUser(this.mongo, tenant, token);
}
// Check to see if this token is a SSO Token or not, if it isn't error out.
@@ -174,7 +173,7 @@ export default class SSOStrategy extends Strategy {
// The token provided does not confirm to the OpenID Connect provider
// spec, but id does conform to a SSOToken so we should expect the token to
// contain the user profile.
return findOrCreateSSOUser(this.db, tenant, token);
return findOrCreateSSOUser(this.mongo, tenant, token);
}
public authenticate(req: Request) {
+9 -5
View File
@@ -1,11 +1,10 @@
import { NextFunction, Response } from "express";
import { Db } from "mongodb";
import { retrieveTenantByDomain } from "talk-server/models/tenant";
import TenantCache from "talk-server/services/tenant/cache";
import { Request } from "talk-server/types/express";
export interface MiddlewareOptions {
db: Db;
cache: TenantCache;
}
export default (options: MiddlewareOptions) => async (
@@ -14,13 +13,18 @@ export default (options: MiddlewareOptions) => async (
next: NextFunction
) => {
try {
// TODO: replace with shared synced cache instead of direct db access.
const tenant = await retrieveTenantByDomain(options.db, req.hostname);
const { cache } = options;
// Attach the tenant to the request.
const tenant = await cache.retrieveByDomain(req.hostname);
if (!tenant) {
// TODO: send a http.StatusNotFound?
return next(new Error("tenant not found"));
}
// Attach the tenant cache to the request.
req.tenantCache = cache;
// Attach the tenant to the request.
req.tenant = tenant;
+7 -2
View File
@@ -33,7 +33,7 @@ async function createTenantRouter(app: AppOptions, options: RouterOptions) {
const router = express.Router();
// Tenant identification middleware.
router.use(tenantMiddleware({ db: app.mongo }));
router.use(tenantMiddleware({ cache: app.tenantCache }));
// Setup Passport middleware.
router.use(options.passport.initialize());
@@ -48,7 +48,12 @@ async function createTenantRouter(app: AppOptions, options: RouterOptions) {
// Any users may submit their GraphQL requests with authentication, this
// middleware will unpack their user into the request.
options.passport.authenticate("jwt", { session: false }),
await tenantGraphMiddleware(app.schemas.tenant, app.config, app.mongo)
await tenantGraphMiddleware({
schema: app.schemas.tenant,
config: app.config,
mongo: app.mongo,
redis: app.redis,
})
);
return router;
-94
View File
@@ -1,94 +0,0 @@
import convict from "convict";
import Joi from "joi";
// Add custom format for the mongo uri scheme.
convict.addFormat({
name: "mongo-uri",
validate: (url: string) => {
Joi.assert(
url,
Joi.string().uri({
scheme: ["mongodb"],
})
);
},
});
// Add custom format for the redis uri scheme.
convict.addFormat({
name: "redis-uri",
validate: (url: string) => {
Joi.assert(
url,
Joi.string().uri({
scheme: ["redis"],
})
);
},
});
const config = convict({
env: {
doc: "The application environment.",
format: ["production", "development", "test"],
default: "development",
env: "NODE_ENV",
},
port: {
doc: "The port to bind.",
format: "port",
default: 3000,
env: "PORT",
arg: "port",
},
mongodb: {
doc: "The MongoDB database to connect to.",
format: "mongo-uri",
default: "mongodb://127.0.0.1:27017/talk",
env: "MONGODB",
arg: "mongodb",
},
redis: {
doc: "The Redis database to connect to.",
format: "redis-uri",
default: "redis://127.0.0.1:6379",
env: "REDIS",
arg: "redis",
},
signing_secret: {
doc: "",
format: "*",
default: "keyboard cat", // TODO: (wyattjoh) evaluate best solution
env: "SIGNING_SECRET",
arg: "signingSecret",
},
signing_algorithm: {
doc: "",
format: [
"HS256",
"HS384",
"HS512",
"RS256",
"RS384",
"RS512",
"ES256",
"ES384",
"ES512",
],
default: "HS256",
env: "SIGNING_ALGORITHM",
arg: "signingAlgorithm",
},
logging_level: {
doc: "The logging level to print to the console",
format: ["fatal", "error", "warn", "info", "debug", "trace"],
default: "info",
env: "LOGGING_LEVEL",
arg: "logging",
},
});
export type Config = typeof config;
// Setup the base configuration.
export default config;
+5 -1
View File
@@ -1,13 +1,17 @@
import { User } from "talk-server/models/user";
import { Request } from "talk-server/types/express";
export interface CommonContextOptions {
user?: User;
req?: Request;
}
export default class CommonContext {
public user?: User;
public req?: Request;
constructor({ user }: CommonContextOptions) {
constructor({ user, req }: CommonContextOptions) {
this.user = user;
this.req = req;
}
}
@@ -5,7 +5,7 @@ import {
GraphQLOptions,
} from "apollo-server-express";
import { FieldDefinitionNode, GraphQLError, ValidationContext } from "graphql";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
// Sourced from: https://github.com/apollographql/apollo-server/blob/958846887598491fadea57b3f9373d129300f250/packages/apollo-server-core/src/ApolloServer.ts#L46-L57
const NoIntrospection = (context: ValidationContext) => ({
@@ -1,5 +1,5 @@
import { RedisPubSub } from "graphql-redis-subscriptions";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
import { createRedisClient } from "talk-server/services/redis";
export async function createPubSub(config: Config): Promise<RedisPubSub> {
+8 -5
View File
@@ -1,16 +1,19 @@
import { Db } from "mongodb";
import CommonContext from "talk-server/graph/common/context";
import { Request } from "talk-server/types/express";
export interface ManagementContextOptions {
db: Db;
mongo: Db;
req?: Request;
}
export default class ManagementContext extends CommonContext {
public db: Db;
public mongo: Db;
constructor({ db }: ManagementContextOptions) {
super({});
constructor({ req, mongo }: ManagementContextOptions) {
super({ req });
this.db = db;
this.mongo = mongo;
}
}
@@ -1,13 +1,14 @@
import { GraphQLSchema } from "graphql";
import { Db } from "mongodb";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
import { graphqlMiddleware } from "talk-server/graph/common/middleware";
import { Request } from "talk-server/types/express";
import Context from "./context";
import ManagementContext from "./context";
export default (schema: GraphQLSchema, config: Config, db: Db) =>
graphqlMiddleware(config, async () => ({
export default (schema: GraphQLSchema, config: Config, mongo: Db) =>
graphqlMiddleware(config, async (req: Request) => ({
schema,
context: new Context({ db }),
context: new ManagementContext({ req, mongo }),
}));
+24 -5
View File
@@ -1,30 +1,49 @@
import { Redis } from "ioredis";
import { Db } from "mongodb";
import CommonContext from "talk-server/graph/common/context";
import { Tenant } from "talk-server/models/tenant";
import { User } from "talk-server/models/user";
import TenantCache from "talk-server/services/tenant/cache";
import { Request } from "talk-server/types/express";
import loaders from "./loaders";
import mutators from "./mutators";
export interface TenantContextOptions {
db: Db;
mongo: Db;
redis: Redis;
tenant: Tenant;
tenantCache: TenantCache;
req?: Request;
user?: User;
}
export default class TenantContext extends CommonContext {
public loaders: ReturnType<typeof loaders>;
public mutators: ReturnType<typeof mutators>;
public db: Db;
public mongo: Db;
public redis: Redis;
public user?: User;
public tenant: Tenant;
public tenantCache: TenantCache;
constructor({ user, tenant, db }: TenantContextOptions) {
super({ user });
constructor({
req,
user,
tenant,
mongo,
redis,
tenantCache,
}: TenantContextOptions) {
super({ user, req });
this.tenant = tenant;
this.tenantCache = tenantCache;
this.user = user;
this.mongo = mongo;
this.redis = redis;
this.loaders = loaders(this);
this.mutators = mutators(this);
this.db = db;
}
}
@@ -10,8 +10,8 @@ import { findOrCreate } from "talk-server/services/assets";
export default (ctx: TenantContext) => ({
findOrCreate: (input: FindOrCreateAssetInput) =>
findOrCreate(ctx.db, ctx.tenant, input),
findOrCreate(ctx.mongo, ctx.tenant, input),
asset: new DataLoader<string, Asset | null>(ids =>
retrieveManyAssets(ctx.db, ctx.tenant.id, ids)
retrieveManyAssets(ctx.mongo, ctx.tenant.id, ids)
),
});
@@ -14,7 +14,7 @@ import {
export default (ctx: Context) => ({
comment: new DataLoader((ids: string[]) =>
retrieveManyComments(ctx.db, ctx.tenant.id, ids)
retrieveManyComments(ctx.mongo, ctx.tenant.id, ids)
),
forAsset: (
assetID: string,
@@ -25,7 +25,7 @@ export default (ctx: Context) => ({
after,
}: AssetToCommentsArgs
) =>
retrieveCommentAssetConnection(ctx.db, ctx.tenant.id, assetID, {
retrieveCommentAssetConnection(ctx.mongo, ctx.tenant.id, assetID, {
first,
orderBy,
after,
@@ -40,9 +40,15 @@ export default (ctx: Context) => ({
after,
}: CommentToRepliesArgs
) =>
retrieveCommentRepliesConnection(ctx.db, ctx.tenant.id, assetID, parentID, {
first,
orderBy,
after,
}),
retrieveCommentRepliesConnection(
ctx.mongo,
ctx.tenant.id,
assetID,
parentID,
{
first,
orderBy,
after,
}
),
});
@@ -4,6 +4,6 @@ import { retrieveManyUsers, User } from "talk-server/models/user";
export default (ctx: Context) => ({
user: new DataLoader<string, User | null>(ids =>
retrieveManyUsers(ctx.db, ctx.tenant.id, ids)
retrieveManyUsers(ctx.mongo, ctx.tenant.id, ids)
),
});
+24 -4
View File
@@ -1,21 +1,41 @@
import { GraphQLSchema } from "graphql";
import { Redis } from "ioredis";
import { Db } from "mongodb";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
import { graphqlMiddleware } from "talk-server/graph/common/middleware";
import { Request } from "talk-server/types/express";
import TenantContext from "./context";
export default async (schema: GraphQLSchema, config: Config, db: Db) => {
export interface TenantGraphQLMiddlewareOptions {
schema: GraphQLSchema;
config: Config;
mongo: Db;
redis: Redis;
}
export default async ({
schema,
config,
mongo,
redis,
}: TenantGraphQLMiddlewareOptions) => {
return graphqlMiddleware(config, async (req: Request) => {
// Load the tenant and user from the request.
const { tenant, user } = req;
const { tenant, user, tenantCache } = req;
// Return the graph options.
return {
schema,
context: new TenantContext({ db, tenant: tenant!, user }),
context: new TenantContext({
req,
mongo,
redis,
tenant: tenant!,
user,
tenantCache,
}),
};
});
};
@@ -5,12 +5,17 @@ import { create } from "talk-server/services/comments";
export default (ctx: TenantContext) => ({
create: (input: GQLCreateCommentInput): Promise<Comment> => {
// FIXME: remove tenant + user !
return create(ctx.db, ctx.tenant, {
author_id: ctx.user!.id,
asset_id: input.assetID,
body: input.body,
parent_id: input.parentID,
});
return create(
ctx.mongo,
ctx.tenant,
ctx.user!,
{
author_id: ctx.user!.id,
asset_id: input.assetID,
body: input.body,
parent_id: input.parentID,
},
ctx.req
);
},
});
@@ -1,6 +1,9 @@
import TenantContext from "talk-server/graph/tenant/context";
import Comment from "./comment";
import Settings from "./settings";
export default (ctx: TenantContext) => ({
Comment: Comment(ctx),
Settings: Settings(ctx),
});
@@ -0,0 +1,11 @@
import { isNull, omitBy } from "lodash";
import TenantContext from "talk-server/graph/tenant/context";
import { GQLSettingsInput } from "talk-server/graph/tenant/schema/__generated__/types";
import { Tenant } from "talk-server/models/tenant";
import { update } from "talk-server/services/tenant";
export default ({ mongo, redis, tenantCache, tenant }: TenantContext) => ({
update: (input: GQLSettingsInput): Promise<Tenant | null> =>
update(mongo, redis, tenantCache, tenant, omitBy(input, isNull)),
});
@@ -1,5 +1,5 @@
import { GQLAuthIntegrationsTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { AuthIntegration, AuthIntegrations } from "talk-server/models/tenant";
import { AuthIntegration, AuthIntegrations } from "talk-server/models/settings";
const disabled: AuthIntegration = { enabled: false };
@@ -1,5 +1,5 @@
import { GQLAuthSettingsTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { Auth } from "talk-server/models/tenant";
import { Auth } from "talk-server/models/settings";
const AuthSettings: GQLAuthSettingsTypeResolver<Auth> = {
integrations: auth => auth.integrations,
@@ -1,5 +1,5 @@
import { GQLFacebookAuthIntegrationTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { FacebookAuthIntegration } from "talk-server/models/tenant";
import { FacebookAuthIntegration } from "talk-server/models/settings";
const FacebookAuthIntegration: GQLFacebookAuthIntegrationTypeResolver<
FacebookAuthIntegration
@@ -1,5 +1,5 @@
import { GQLGoogleAuthIntegrationTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { GoogleAuthIntegration } from "talk-server/models/tenant";
import { GoogleAuthIntegration } from "talk-server/models/settings";
const GoogleAuthIntegration: GQLGoogleAuthIntegrationTypeResolver<
GoogleAuthIntegration
@@ -1,5 +1,5 @@
import { GQLLocalAuthIntegrationTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { LocalAuthIntegration } from "talk-server/models/tenant";
import { LocalAuthIntegration } from "talk-server/models/settings";
const LocalAuthIntegration: GQLLocalAuthIntegrationTypeResolver<
LocalAuthIntegration
@@ -5,6 +5,10 @@ const Mutation: GQLMutationTypeResolver<void> = {
comment: await ctx.mutators.Comment.create(input),
clientMutationId: input.clientMutationId,
}),
updateSettings: async (source, { input }, ctx) => ({
settings: await ctx.mutators.Settings.update(input.settings),
clientMutationId: input.clientMutationId,
}),
};
export default Mutation;
@@ -1,5 +1,5 @@
import { GQLOIDCAuthIntegrationTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { OIDCAuthIntegration } from "talk-server/models/tenant";
import { OIDCAuthIntegration } from "talk-server/models/settings";
const OIDCAuthIntegration: GQLOIDCAuthIntegrationTypeResolver<
OIDCAuthIntegration
@@ -1,5 +1,5 @@
import { GQLSSOAuthIntegrationTypeResolver } from "talk-server/graph/tenant/schema/__generated__/types";
import { SSOAuthIntegration } from "talk-server/models/tenant";
import { SSOAuthIntegration } from "talk-server/models/settings";
const SSOAuthIntegration: GQLSSOAuthIntegrationTypeResolver<
SSOAuthIntegration
@@ -25,6 +25,19 @@ Cursor represents a paginating cursor.
"""
scalar Cursor
################################################################################
## Actions
################################################################################
enum ACTION_TYPE {
FLAG
DONTAGREE
}
enum ACTION_ITEM_TYPE {
COMMENTS
}
################################################################################
## Settings
################################################################################
@@ -287,7 +300,7 @@ type Settings {
"""
domains will return a given list of whitelisted domains.
"""
domains: [String!] @auth(roles: [ADMIN]) @auth(roles: [ADMIN])
domains: [String!] @auth(roles: [ADMIN])
"""
auth contains all the settings related to authentication and authorization.
@@ -301,6 +314,7 @@ type Settings {
enum USER_ROLE {
COMMENTER
STAFF
MODERATOR
ADMIN
}
@@ -390,8 +404,34 @@ type User {
################################################################################
enum COMMENT_STATUS {
"""
The comment is not PREMOD, but was not applied a moderation status by a
moderator.
"""
NONE
"""
The comment has been accepted by a moderator.
"""
ACCEPTED
"""
The comment has been rejected by a moderator.
"""
REJECTED
"""
The comment was created while the asset's premoderation option was on, and
new comments that haven't been moderated yet are referred to as
"premoderated" or "premod" comments.
"""
PREMOD
"""
SYSTEM_WITHHELD represents a comment that was withheld by the system because
it was flagged by an internal process for further review.
"""
SYSTEM_WITHHELD
}
"""
@@ -661,6 +701,66 @@ type CreateCommentPayload {
clientMutationId: String!
}
##################
## updateSettings
##################
"""
SettingsInput is the partial type of the Settings type for performing mutations.
"""
input SettingsInput {
moderation: MODERATION_MODE
requireEmailConfirmation: Boolean
infoBoxEnable: Boolean
infoBoxContent: String
questionBoxEnable: Boolean
questionBoxContent: String
questionBoxIcon: String
premodLinksEnable: Boolean
autoCloseStream: Boolean
customCssUrl: String
closedTimeout: Int
closedMessage: String
disableCommenting: Boolean
disableCommentingMessage: String
editCommentWindowLength: Int
charCountEnable: Boolean
charCount: Int
organizationName: String
organizationContactEmail: String
# wordlist: WordlistSettings @auth(roles: [ADMIN, MODERATOR])
domains: [String!]
# auth: AuthSettings!
}
"""
UpdateSettingsInput provides the input for the updateSettings Mutation.
"""
input UpdateSettingsInput {
settings: SettingsInput!
"""
clientMutationId is required for Relay support.
"""
clientMutationId: String!
}
"""
UpdateSettingsPayload contains the updated Settings after the updateSettings
mutation.
"""
type UpdateSettingsPayload {
"""
settings is the updated Settings.
"""
settings: Settings
"""
clientMutationId is required for Relay support.
"""
clientMutationId: String!
}
##################
## Mutation
##################
@@ -670,6 +770,11 @@ type Mutation {
createComment will create a Comment as the current logged in User.
"""
createComment(input: CreateCommentInput!): CreateCommentPayload @auth
"""
updateSettings will update the Settings for the given Tenant.
"""
updateSettings(input: UpdateSettingsInput!): UpdateSettingsPayload @auth(roles: [ADMIN])
}
################################################################################
+9 -1
View File
@@ -1,13 +1,14 @@
import express, { Express } from "express";
import http from "http";
import config, { Config } from "talk-common/config";
import { createJWTSigningConfig } from "talk-server/app/middleware/passport/jwt";
import getManagementSchema from "talk-server/graph/management/schema";
import { Schemas } from "talk-server/graph/schemas";
import getTenantSchema from "talk-server/graph/tenant/schema";
import TenantCache from "talk-server/services/tenant/cache";
import { attachSubscriptionHandlers, createApp, listenAndServe } from "./app";
import config, { Config } from "./config";
import logger from "./logger";
import { createMongoDB } from "./services/mongodb";
import { createRedisClient } from "./services/redis";
@@ -68,6 +69,12 @@ class Server {
// Create the signing config.
const signingConfig = createJWTSigningConfig(this.config);
// Create the TenantCache.
const tenantCache = new TenantCache(mongo, await createRedisClient(config));
// Prime the tenant cache so it'll be ready to serve now.
await tenantCache.primeAll();
// Create the Talk App, branching off from the parent app.
const app: Express = await createApp({
parent,
@@ -76,6 +83,7 @@ class Server {
config: this.config,
schemas: this.schemas,
signingConfig,
tenantCache,
});
// Start the application and store the resulting http.Server.
+5 -1
View File
@@ -1,8 +1,12 @@
import bunyan from "bunyan";
import bunyan, { LogLevelString } from "bunyan";
import config from "talk-common/config";
const logger = bunyan.createLogger({
name: "talk",
serializers: bunyan.stdSerializers,
// TODO: (wyattjoh) move this into some managed instance?
level: config.get("logging_level") as LogLevelString,
});
export default logger;
+16 -2
View File
@@ -1,3 +1,17 @@
export interface ActionCounts {
[_: string]: number;
import {
GQLACTION_ITEM_TYPE,
GQLACTION_TYPE,
} from "talk-server/graph/tenant/schema/__generated__/types";
export type ActionCounts = Record<string, number>;
export interface Action {
readonly id: string;
action_type: GQLACTION_TYPE;
item_type: GQLACTION_ITEM_TYPE;
item_id: string;
group_id?: string;
user_id?: string;
created_at: Date;
metadata?: Record<string, any>;
}
+8 -1
View File
@@ -4,6 +4,7 @@ import { Db } from "mongodb";
import uuid from "uuid";
import { Omit } from "talk-common/types";
import { ModerationSettings } from "talk-server/models/settings";
import { TenantResource } from "talk-server/models/tenant";
function collection(db: Db) {
@@ -25,6 +26,12 @@ export interface Asset extends TenantResource {
publication_date?: Date;
modified_date?: Date;
created_at: Date;
/**
* settings provides a point where the settings can be overriden for a
* specific Asset.
*/
settings?: Partial<ModerationSettings>;
}
export interface UpsertAssetInput {
@@ -170,7 +177,7 @@ export async function updateAsset(
const result = await collection(db).findOneAndUpdate(
{ id, tenant_id: tenantID },
// Only update fields that have been updated.
{ $set: dotize(update) },
{ $set: dotize.convert(update) },
// False to return the updated document instead of the original
// document.
{ returnOriginal: false }
+4 -8
View File
@@ -1,4 +1,3 @@
import { merge } from "lodash";
import { Db } from "mongodb";
import uuid from "uuid";
@@ -91,17 +90,14 @@ export async function createComment(
};
// Merge the defaults and the input together.
const comment: Readonly<Comment> = merge({}, defaults, input);
// TODO: Check for existence of the parent ID before we create the comment.
// TODO: Check for existence of the asset ID before we create the comment.
const comment: Readonly<Comment> = {
...defaults,
...input,
};
// Insert it into the database.
await collection(db).insertOne(comment);
// TODO: update reply count of parent if exists.
return comment;
}
+246
View File
@@ -0,0 +1,246 @@
import {
GQLMODERATION_MODE,
GQLUSER_ROLE,
} from "talk-server/graph/tenant/schema/__generated__/types";
export interface Wordlist {
banned: string[];
suspect: string[];
}
export interface EmailDomainRuleCondition {
/**
* emailDomain is the domain name component of the email addresses that should
* match for this condition.
*/
emailDomain: string;
/**
* emailVerifiedRequired stipulates that this rule only applies when the user
* account has been marked as having their email address already verified.
*/
emailVerifiedRequired: boolean;
}
/**
* RoleRule describes the role assignment for when a user logs into Talk, how
* they can have their account automatically upgraded to a specific role when
* the domain for their email matches the one provided.
*/
export interface RoleRule extends Partial<EmailDomainRuleCondition> {
/**
* role is the specific GQLUSER_ROLE that should be assigned to the newly
* created user depending on their email address.
*/
role: GQLUSER_ROLE;
}
export interface AuthRules {
/**
* roles allow the configuration of automatic role assignment based on the
* user's email address.
*/
roles?: RoleRule[];
/**
* restrictTo when populated, will restrict which users can login using this
* integration. If a user successfully logs in using the OIDCStrategy, but
* does not match the following rules, the user will not be created.
*/
restrictTo?: EmailDomainRuleCondition[];
}
export interface AuthIntegration {
enabled: boolean;
}
export interface DisplayNameAuthIntegration {
displayNameEnable: boolean;
}
/**
* SSOAuthIntegration is an AuthIntegration that provides a secret to the admins
* of a tenant, where they can sign a SSO payload with it to provide to the
* embed to allow single sign on.
*/
export interface SSOAuthIntegration
extends AuthIntegration,
DisplayNameAuthIntegration {
key: string;
}
/**
* OIDCAuthIntegration provides a way to store Open ID Connect credentials. This
* will be used in the admin to provide staff logins for users.
*/
export interface OIDCAuthIntegration
extends AuthIntegration,
DisplayNameAuthIntegration {
clientID: string;
clientSecret: string;
issuer: string;
authorizationURL: string;
jwksURI: string;
tokenURL: string;
}
export interface FacebookAuthIntegration extends AuthIntegration {
clientID: string;
clientSecret: string;
}
export interface GoogleAuthIntegration extends AuthIntegration {
clientID: string;
clientSecret: string;
}
export type LocalAuthIntegration = AuthIntegration;
/**
* AuthIntegrations describes all of the possible auth integration
* configurations.
*/
export interface AuthIntegrations {
/**
* local is the auth integration for the email/password based auth.
*/
local: LocalAuthIntegration;
/**
* sso is the external auth integration for the single sign on auth.
*/
sso?: SSOAuthIntegration;
/**
* sso is the external auth integration for the OpenID Connect auth.
*/
oidc?: OIDCAuthIntegration;
/**
* sso is the external auth integration for the Google auth.
*/
google?: GoogleAuthIntegration;
/**
* sso is the external auth integration for the Facebook auth.
*/
facebook?: FacebookAuthIntegration;
}
export interface Auth {
integrations: AuthIntegrations;
}
/**
* Akismet provides integration with the Akismet Spam detection service.
*/
export interface AkismetIntegration {
/**
* When true, it will enable comments to be checked by Akismet.
*/
enabled: boolean;
/**
* The key for the Akismet integration.
*/
key?: string;
/**
* The site (blog) for the Akismet integration.
*/
site?: string;
}
export interface ExternalIntegrations {
/**
* akismet provides integration with the Akismet Spam detection service.
*/
akismet: AkismetIntegration;
}
export interface ModerationSettings {
moderation: GQLMODERATION_MODE;
requireEmailConfirmation: boolean;
infoBoxEnable: boolean;
infoBoxContent?: string;
questionBoxEnable: boolean;
questionBoxIcon?: string;
questionBoxContent?: string;
premodLinksEnable: boolean;
autoCloseStream: boolean;
closedTimeout: number;
closedMessage?: string;
disableCommenting: boolean;
disableCommentingMessage?: string;
charCountEnable: boolean;
charCount?: number;
}
/**
* KarmaThreshold defines the bounds for which a User will become unreliable or
* reliable based on their karma score. If the score is equal or less than the
* unreliable value, they are unreliable. If the score is equal or more than the
* reliable value, they are reliable. If they are neither reliable or unreliable
* then they are neutral.
*/
export interface KarmaThreshold {
reliable: number;
unreliable: number;
}
export interface KarmaThresholds {
/**
* flag represents karma settings in relation to how well a User's flagging
* ability aligns with the moderation decicions made by moderators.
*/
flag: KarmaThreshold;
/**
* comment represents the karma setting in relation to how well a User's comments are moderated.
*/
comment: KarmaThreshold;
}
export interface Karma {
/**
* When true, checks will be completed to ensure that the Karma checks are
* completed.
*/
enabled: boolean;
/**
* karmaThresholds contains the currently set thresholds for triggering Trust
* beheviour.
*/
thresholds: KarmaThresholds;
}
export interface Settings extends ModerationSettings {
customCssUrl?: string;
/**
* editCommentWindowLength is the length of time (in milliseconds) after a
* comment is posted that it can still be edited by the author.
*/
editCommentWindowLength: number;
/**
* karma is the set of settings related to how user Trust and Karma are
* handled.
*/
karma: Karma;
/**
* wordlist stores all the banned/suspect words.
*/
wordlist: Wordlist;
/**
* Set of configured authentication integrations.
*/
auth: Auth;
/**
* Various integrations with external services.
*/
integrations: ExternalIntegrations;
}
+32 -142
View File
@@ -1,13 +1,10 @@
import dotize from "dotize";
import { merge } from "lodash";
import { Db } from "mongodb";
import uuid from "uuid";
import { Sub } from "talk-common/types";
import {
GQLMODERATION_MODE,
GQLUSER_ROLE,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { Omit, Sub } from "talk-common/types";
import { GQLMODERATION_MODE } from "talk-server/graph/tenant/schema/__generated__/types";
import { Settings } from "talk-server/models/settings";
function collection(db: Db) {
return db.collection<Readonly<Tenant>>("tenants");
@@ -17,147 +14,21 @@ export interface TenantResource {
readonly tenant_id: string;
}
export interface Wordlist {
banned: string[];
suspect: string[];
}
// AuthIntegrations.
export interface EmailDomainRuleCondition {
// emailDomain is the domain name component of the email addresses that should
// match for this condition.
emailDomain: string;
// emailVerifiedRequired stipulates that this rule only applies when the user
// account has been marked as having their email address already verified.
emailVerifiedRequired: boolean;
}
// RoleRule describes the role assignment for when a user logs into Talk, how
// they can have their account automatically upgraded to a specific role when
// the domain for their email matches the one provided.
export interface RoleRule extends Partial<EmailDomainRuleCondition> {
// role is the specific GQLUSER_ROLE that should be assigned to the newly created
// user depending on their email address.
role: GQLUSER_ROLE;
}
export interface AuthRules {
// roles allow the configuration of automatic role assignment based on the
// user's email address.
roles?: RoleRule[];
// restrictTo when populated, will restrict which users can login using this
// integration. If a user successfully logs in using the OIDCStrategy, but
// does not match the following rules, the user will not be created.
restrictTo?: EmailDomainRuleCondition[];
}
export interface AuthIntegration {
enabled: boolean;
}
export interface DisplayNameAuthIntegration {
displayNameEnable: boolean;
}
// SSOAuthIntegration is an AuthIntegration that provides a secret to the admins
// of a tenant, where they can sign a SSO payload with it to provide to the
// embed to allow single sign on.
export interface SSOAuthIntegration
extends AuthIntegration,
DisplayNameAuthIntegration {
key: string;
}
// OIDCAuthIntegration provides a way to store Open ID Connect credentials. This
// will be used in the admin to provide staff logins for users.
export interface OIDCAuthIntegration
extends AuthIntegration,
DisplayNameAuthIntegration {
clientID: string;
clientSecret: string;
issuer: string;
authorizationURL: string;
jwksURI: string;
tokenURL: string;
}
export interface FacebookAuthIntegration extends AuthIntegration {
clientID: string;
clientSecret: string;
}
export interface GoogleAuthIntegration extends AuthIntegration {
clientID: string;
clientSecret: string;
}
export type LocalAuthIntegration = AuthIntegration;
// AuthIntegrations describes all of the possible auth integration configurations.
export interface AuthIntegrations {
// local is the auth integration for the local auth.
local: LocalAuthIntegration;
// sso is the external auth integration for the single sign on auth.
sso?: SSOAuthIntegration;
// sso is the external auth integration for the OpenID Connect auth.
oidc?: OIDCAuthIntegration;
// sso is the external auth integration for the Google auth.
google?: GoogleAuthIntegration;
// sso is the external auth integration for the Facebook auth.
facebook?: FacebookAuthIntegration;
}
export interface Auth {
integrations: AuthIntegrations;
}
// Tenant definition.
export interface Tenant {
/**
* Tenant describes a given Tenant on Talk that has Assets, Comments, and Users.
*/
export interface Tenant extends Settings {
readonly id: string;
// Domain is set when the tenant is created, and is used to retrieve the
// specific tenant that the API request pertains to.
domain: string;
moderation: GQLMODERATION_MODE;
requireEmailConfirmation: boolean;
infoBoxEnable: boolean;
infoBoxContent?: string;
questionBoxEnable: boolean;
questionBoxIcon?: string;
questionBoxContent?: string;
premodLinksEnable: boolean;
autoCloseStream: boolean;
closedTimeout: number;
closedMessage?: string;
customCssUrl?: string;
disableCommenting: boolean;
disableCommentingMessage?: string;
// editCommentWindowLength is the length of time (in milliseconds) after a
// comment is posted that it can still be edited by the author.
editCommentWindowLength: number;
charCountEnable: boolean;
charCount?: number;
organizationName: string;
organizationContactEmail: string;
// wordlist stores all the banned/suspect words.
wordlist: Wordlist;
// domains is the set of whitelisted domains.
// domains is the list of domains that are allowed to have the iframe load on.
domains: string[];
// Set of configured authentication integrations.
auth: Auth;
organizationName: string;
organizationContactEmail: string;
}
/**
@@ -207,10 +78,27 @@ export async function createTenant(db: Db, input: CreateTenantInput) {
},
},
},
karma: {
enabled: true,
thresholds: {
// By default, flaggers are reliable after one correct flag, and
// unreliable if there is an incorrect flag.
flag: { reliable: 1, unreliable: -1 },
comment: { reliable: 1, unreliable: -1 },
},
},
integrations: {
akismet: {
enabled: false,
},
},
};
// Create the new Tenant by merging it together with the defaults.
const tenant: Readonly<Tenant> = merge({}, input, defaults);
const tenant: Readonly<Tenant> = {
...defaults,
...input,
};
// Insert the Tenant into the database.
await collection(db).insert(tenant);
@@ -258,16 +146,18 @@ export async function retrieveAllTenants(db: Db) {
.toArray();
}
export type UpdateTenantInput = Omit<Partial<Tenant>, "id" | "domain">;
export async function updateTenant(
db: Db,
id: string,
update: Partial<CreateTenantInput>
update: UpdateTenantInput
) {
// Get the tenant from the database.
const result = await collection(db).findOneAndUpdate(
{ id },
// Only update fields that have been updated.
{ $set: dotize(update) },
{ $set: dotize.convert(update) },
// False to return the updated document instead of the original
// document.
{ returnOriginal: false }
+5 -3
View File
@@ -1,5 +1,4 @@
import bcrypt from "bcryptjs";
import { merge } from "lodash";
import { Db } from "mongodb";
import uuid from "uuid";
@@ -131,11 +130,14 @@ export async function upsertUser(
}
// Merge the defaults and the input together.
const user: Readonly<User> = merge({}, defaults, input, {
const user: Readonly<User> = {
...defaults,
...input,
// Specified last in the merge call, it will override any existing password
// entry if it is defined.
password: hashedPassword,
});
};
// Create a query that will utilize a findOneAndUpdate to facilitate an upsert
// operation to ensure no user has the same profile and/or email address. If
+52 -6
View File
@@ -1,22 +1,68 @@
import { Db } from "mongodb";
import { Omit } from "talk-common/types";
import { GQLCOMMENT_STATUS } from "talk-server/graph/tenant/schema/__generated__/types";
import { createComment, CreateCommentInput } from "talk-server/models/comment";
import { retrieveAsset } from "talk-server/models/asset";
import {
createComment,
CreateCommentInput,
retrieveComment,
} from "talk-server/models/comment";
import { Tenant } from "talk-server/models/tenant";
import { User } from "talk-server/models/user";
import { processForModeration } from "talk-server/services/comments/moderation";
import { Request } from "talk-server/types/express";
export type CreateComment = Omit<
CreateCommentInput,
"status" | "action_counts"
>;
export async function create(db: Db, tenant: Tenant, input: CreateComment) {
// TODO: run the comment through the moderation phases.
const comment = await createComment(db, tenant.id, {
status: GQLCOMMENT_STATUS.ACCEPTED,
export async function create(
mongo: Db,
tenant: Tenant,
author: User,
input: CreateComment,
req?: Request
) {
const asset = await retrieveAsset(mongo, tenant.id, input.asset_id);
if (!asset) {
// TODO: (wyattjoh) return better error.
throw new Error("asset referenced does not exist");
}
// TODO: (wyattjoh) Check that the asset was visable.
if (input.parent_id) {
// Check to see that the reference parent ID exists.
const parent = await retrieveComment(mongo, tenant.id, input.parent_id);
if (!parent) {
// TODO: (wyattjoh) return better error.
throw new Error("parent comment referenced does not exist");
}
// TODO: (wyattjoh) Check that the parent comment was visible.
}
// Run the comment through the moderation phases.
const { status } = await processForModeration({
asset,
tenant,
comment: input,
author,
req,
});
// TODO: (wyattjoh) use the actions somehow.
const comment = await createComment(mongo, tenant.id, {
status,
action_counts: {},
...input,
});
if (input.parent_id) {
// TODO: update reply count of parent.
}
return comment;
}
@@ -0,0 +1,75 @@
import { Omit, Promiseable } from "talk-common/types";
import { GQLCOMMENT_STATUS } from "talk-server/graph/tenant/schema/__generated__/types";
import { Action } from "talk-server/models/actions";
import { Asset } from "talk-server/models/asset";
import { Tenant } from "talk-server/models/tenant";
import { CreateComment } from "talk-server/services/comments";
import { User } from "talk-server/models/user";
import { Request } from "talk-server/types/express";
import { moderationPhases } from "./phases";
// TODO: (wyattjoh) move into actions module.
export type CreateAction = Omit<
Action,
"id" | "item_type" | "item_id" | "created_at"
>;
export interface PhaseResult {
actions: CreateAction[];
status: GQLCOMMENT_STATUS;
}
export interface ModerationPhaseContext {
asset: Asset;
tenant: Tenant;
comment: CreateComment;
author: User;
req?: Request;
}
export type ModerationPhase = (
context: ModerationPhaseContext
) => Promiseable<PhaseResult>;
export type IntermediatePhaseResult = Partial<PhaseResult> | void;
export type IntermediateModerationPhase = (
context: ModerationPhaseContext
) => Promiseable<IntermediatePhaseResult>;
/**
* compose will create a moderation pipeline for which is executable with the
* passed actions.
*/
const compose = (
phases: IntermediateModerationPhase[]
): ModerationPhase => async context => {
const actions: CreateAction[] = [];
// Loop over all the moderation phases and see if we've resolved the status.
for (const phase of phases) {
const result = await phase(context);
if (result) {
if (result.actions) {
actions.push(...result.actions);
}
// If this result contained a status, then we've finished resolving
// phases!
const { status } = result;
if (status) {
return { status, actions };
}
}
}
// If we didn't determine a different comment from a previous itteration, set
// it to 'NONE'.
return { status: GQLCOMMENT_STATUS.NONE, actions };
};
/**
* process the comment and return moderation details.
*/
export const processForModeration: ModerationPhase = compose(moderationPhases);
@@ -0,0 +1,42 @@
import { Asset } from "talk-server/models/asset";
import { Comment } from "talk-server/models/comment";
import { Tenant } from "talk-server/models/tenant";
import { User } from "talk-server/models/user";
import { assetClosed } from "talk-server/services/comments/moderation/phases/assetClosed";
describe("assetClosed", () => {
it("throws an error when the asset is closed", () => {
const asset = { closedAt: new Date() };
expect(() =>
assetClosed({
asset: asset as Asset,
tenant: (null as any) as Tenant,
comment: (null as any) as Comment,
author: (null as any) as User,
})
).toThrow();
});
it("does not throw an error when the asset is not closed", () => {
const now = new Date();
expect(
assetClosed({
asset: { closedAt: new Date(now.getTime() + 60000) } as Asset,
tenant: (null as any) as Tenant,
comment: (null as any) as Comment,
author: (null as any) as User,
})
).toBeUndefined();
expect(
assetClosed({
asset: {} as Asset,
tenant: (null as any) as Tenant,
comment: (null as any) as Comment,
author: (null as any) as User,
})
).toBeUndefined();
});
});
@@ -0,0 +1,12 @@
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
// This phase checks to see if the asset being processed is closed or not.
export const assetClosed: IntermediateModerationPhase = ({ asset }) => {
// Check to see if the asset has closed commenting...
if (asset.closedAt && asset.closedAt.valueOf() <= Date.now()) {
// TODO: (wyattjoh) return better error.
throw new Error("asset is currently closed for commenting");
}
return;
};
@@ -0,0 +1,45 @@
import {
GQLACTION_TYPE,
GQLCOMMENT_STATUS,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { ModerationSettings } from "talk-server/models/settings";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
const testCharCount = (settings: Partial<ModerationSettings>, length: number) =>
settings.charCountEnable && settings.charCount && length > settings.charCount;
export const commentLength: IntermediateModerationPhase = async ({
asset,
tenant,
comment,
}) => {
const length = comment.body.length;
// Check to see if the body is too short, if it is, then complain about it!
if (length < 2) {
// TODO: (wyattjoh) return better error.
throw new Error("comment body too short");
}
// Reject if the comment is too long
if (
testCharCount(tenant, length) ||
(asset.settings && testCharCount(asset.settings, length))
) {
// Add the flag related to Trust to the comment.
return {
status: GQLCOMMENT_STATUS.REJECTED,
actions: [
{
action_type: GQLACTION_TYPE.FLAG,
group_id: "BODY_COUNT",
metadata: {
count: length,
},
},
],
};
}
return;
};
@@ -0,0 +1,19 @@
import { ModerationSettings } from "talk-server/models/settings";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
const testDisabledCommenting = (settings: Partial<ModerationSettings>) =>
settings.disableCommenting;
export const commentingDisabled: IntermediateModerationPhase = ({
asset,
tenant,
}) => {
// Check to see if the asset has closed commenting.
if (
testDisabledCommenting(tenant) ||
(asset.settings && testDisabledCommenting(asset.settings))
) {
// TODO: (wyattjoh) return better error.
throw new Error("commenting has been disabled tenant wide");
}
};
@@ -0,0 +1,26 @@
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
import { premod } from "talk-server/services/comments/moderation/phases/premod";
import { assetClosed } from "./assetClosed";
import { commentingDisabled } from "./commentingDisabled";
import { commentLength } from "./commentLength";
import { karma } from "./karma";
import { links } from "./links";
import { spam } from "./spam";
import { staff } from "./staff";
import { wordlist } from "./wordlist";
/**
* The moderation phases to apply for each comment being processed.
*/
export const moderationPhases: IntermediateModerationPhase[] = [
commentLength,
assetClosed,
commentingDisabled,
wordlist,
staff,
links,
karma,
spam,
premod,
];
@@ -0,0 +1,40 @@
import {
GQLACTION_TYPE,
GQLCOMMENT_STATUS,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
import {
getCommentTrustScore,
isReliableCommenter,
} from "talk-server/services/users/karma";
// This phase checks to see if the user making the comment is allowed to do so
// considering their reliability (Trust) status.
export const karma: IntermediateModerationPhase = ({ tenant, author }) => {
// If the user is not a reliable commenter (passed the unreliability
// threshold by having too many rejected comments) then we can change the
// status of the comment to `SYSTEM_WITHHELD`, therefore pushing the user's
// comments away from the public eye until a moderator can manage them. This
// of course can only be applied if the comment's current status is `NONE`,
// we don't want to interfere if the comment was rejected.
if (
tenant.karma.enabled &&
isReliableCommenter(tenant.karma.thresholds, author) === false
) {
// Add the flag related to Trust to the comment.
return {
status: GQLCOMMENT_STATUS.SYSTEM_WITHHELD,
actions: [
{
action_type: GQLACTION_TYPE.FLAG,
group_id: "TRUST",
metadata: {
trust: getCommentTrustScore(author),
},
},
],
};
}
return;
};
@@ -0,0 +1,49 @@
import linkify from "linkify-it";
import tlds from "tlds";
import {
GQLACTION_TYPE,
GQLCOMMENT_STATUS,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { ModerationSettings } from "talk-server/models/settings";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
/**
* The preloaded linkify instance with common tlds.
*/
const testForLinks = linkify().tlds(tlds);
const testPremodLinksEnable = (
settings: Partial<ModerationSettings>,
body: string
) => settings.premodLinksEnable && testForLinks.test(body);
// This phase checks the comment if it has any links in it if the check is
// enabled.
export const links: IntermediateModerationPhase = ({
asset,
tenant,
comment,
author,
}) => {
if (
testPremodLinksEnable(tenant, comment.body) ||
(asset.settings && testPremodLinksEnable(asset.settings, comment.body))
) {
// Add the flag related to Trust to the comment.
return {
status: GQLCOMMENT_STATUS.SYSTEM_WITHHELD,
actions: [
{
action_type: GQLACTION_TYPE.FLAG,
group_id: "LINKS",
metadata: {
links: comment.body,
},
},
],
};
}
return;
};
@@ -0,0 +1,28 @@
import {
GQLCOMMENT_STATUS,
GQLMODERATION_MODE,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { ModerationSettings } from "talk-server/models/settings";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
const testModerationMode = (settings: Partial<ModerationSettings>) =>
settings.moderation === GQLMODERATION_MODE.PRE;
// This phase checks to see if the settings have premod enabled, if they do,
// the comment is premod, otherwise, it's just none.
export const premod: IntermediateModerationPhase = ({ asset, tenant }) => {
// If the settings say that we're in premod mode, then the comment is in
// premod status.
// TODO: (wyattjoh) pull from the asset settings.
if (
testModerationMode(tenant) ||
(asset.settings && testModerationMode(asset.settings))
) {
return {
status: GQLCOMMENT_STATUS.PREMOD,
};
}
return;
};
@@ -0,0 +1,74 @@
import { Client } from "akismet-api";
import {
GQLACTION_TYPE,
GQLCOMMENT_STATUS,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
export const spam: IntermediateModerationPhase = async ({
asset,
tenant,
comment,
author,
req,
}) => {
const integration = tenant.integrations.akismet;
// We can only check for spam if this comment originated from a graphql
// request via an HTTP call.
if (!req || !integration.enabled) {
return;
}
if (!integration.key || !integration.site) {
return;
}
// Create the Akismet client.
const client = new Client({
key: integration.key,
blog: integration.site,
});
// Grab the properties we need.
const userIP = req.ip;
if (!userIP) {
return;
}
const userAgent = req.get("User-Agent");
if (!userAgent || userAgent.length === 0) {
return;
}
const referrer = req.get("Referrer");
if (!referrer || referrer.length === 0) {
return;
}
// Check the comment for spam.
const isSpam = await client.checkSpam({
user_ip: userIP, // REQUIRED
referrer, // REQUIRED
user_agent: userAgent, // REQUIRED
comment_content: comment.body,
permalink: asset.url,
comment_author: author.displayName || author.username || "",
comment_type: "comment",
is_test: false,
});
if (isSpam) {
return {
status: GQLCOMMENT_STATUS.SYSTEM_WITHHELD,
actions: [
{
action_type: GQLACTION_TYPE.FLAG,
group_id: "SPAM_COMMENT",
},
],
};
}
return;
};
@@ -0,0 +1,21 @@
import {
GQLCOMMENT_STATUS,
GQLUSER_ROLE,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
// If a given user is a staff member, always approve their comment.
export const staff: IntermediateModerationPhase = ({
asset,
tenant,
comment,
author,
}) => {
if (author.role !== GQLUSER_ROLE.COMMENTER) {
return {
status: GQLCOMMENT_STATUS.ACCEPTED,
};
}
return;
};
@@ -0,0 +1,50 @@
import {
GQLACTION_TYPE,
GQLCOMMENT_STATUS,
} from "talk-server/graph/tenant/schema/__generated__/types";
import { IntermediateModerationPhase } from "talk-server/services/comments/moderation";
import { containsMatchingPhrase } from "talk-server/services/comments/moderation/wordlist";
// This phase checks the comment against the wordlist.
export const wordlist: IntermediateModerationPhase = ({
asset,
tenant,
comment,
author,
}) => {
// Decide the status based on whether or not the current asset/settings
// has pre-mod enabled or not. If the comment was rejected based on the
// wordlist, then reject it, otherwise if the moderation setting is
// premod, set it to `premod`.
if (containsMatchingPhrase(tenant.wordlist.banned, comment.body)) {
// Add the flag related to Trust to the comment.
return {
status: GQLCOMMENT_STATUS.REJECTED,
actions: [
{
action_type: GQLACTION_TYPE.FLAG,
group_id: "BANNED_WORD",
},
],
};
}
// If the comment has a suspect word or a link, we need to add a
// flag to it to indicate that it needs to be looked at.
// Otherwise just return the new comment.
// If the wordlist has matched the suspect word filter and we haven't disabled
// auto-flagging suspect words, then we should flag the comment!
if (containsMatchingPhrase(tenant.wordlist.suspect, comment.body)) {
return {
actions: [
{
action_type: GQLACTION_TYPE.FLAG,
group_id: "SUSPECT_WORD",
},
],
};
}
return;
};
@@ -0,0 +1,46 @@
import { containsMatchingPhrase } from "talk-server/services/comments/moderation/wordlist";
const phrases = [
"cookies",
"how to do bad things",
"how to do really bad things",
"s h i t",
"$hit",
"p**ch",
"p*ch",
];
describe("containsMatchingPhrase", () => {
it("does match on a word in the list", () => {
[
"how to do really bad things",
"what is cookies",
"cookies",
"COOKIES.",
"how to do bad things",
"How To do bad things!",
"This stuff is $hit!",
"That's a p**ch!",
].forEach(word => {
expect(containsMatchingPhrase(phrases, word)).toEqual(true);
});
});
it("does not match on a word not in the list", () => {
[
"how to",
"cookie",
"how to be a great person?",
"how to not do really bad things?",
"i have $100 dollars.",
"I have bad $ hit lling",
"That's a p***ch!",
].forEach(word => {
expect(containsMatchingPhrase(phrases, word)).toEqual(false);
});
});
it("allows an empty list", () => {
expect(containsMatchingPhrase([], "test")).toEqual(false);
});
});
@@ -0,0 +1,25 @@
/**
* Escape string for special regular expression characters.
*/
export function escapeRegExp(str: string) {
return str.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); // $& means the whole matched string
}
/**
* Generate a regular expression that catches the `phrases`.
*/
export function generateRegExp(phrases: string[]) {
const inner = phrases
.map(phrase =>
phrase
.split(/\s+/)
.map(word => escapeRegExp(word))
.join('[\\s"?!.]+')
)
.join("|");
return new RegExp(`(^|[^\\w])(${inner})(?=[^\\w]|$)`, "iu");
}
export const containsMatchingPhrase = (phrases: string[], testString: string) =>
phrases.length > 0 ? generateRegExp(phrases).test(testString) : false;
+1 -1
View File
@@ -1,5 +1,5 @@
import { Db, MongoClient } from "mongodb";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
/**
* create will connect to the MongoDB instance identified in the configuration.
+1 -1
View File
@@ -1,5 +1,5 @@
import RedisClient, { Redis } from "ioredis";
import { Config } from "talk-server/config";
import { Config } from "talk-common/config";
/**
* create will connect to the Redis instance identified in the configuration.
+117 -27
View File
@@ -1,35 +1,70 @@
import DataLoader from "dataloader";
import { Redis } from "ioredis";
import { Db } from "mongodb";
import uuid from "uuid";
import { EventEmitter } from "events";
import logger from "talk-server/logger";
import {
retrieveAllTenants,
retrieveManyTenants,
retrieveManyTenantsByDomain,
Tenant,
} from "talk-server/models/tenant";
const CacheUpdateChannel = "tenant";
const TENANT_UPDATE_CHANNEL = "tenant";
// Cache provides an interface for retrieving tenant stored in local memory
// rather than grabbing it from the database every single call.
export default class Cache {
// private tenants: Map<string, Promise<Readonly<Tenant>>>;
private tenants: DataLoader<string, Readonly<Tenant> | null>;
private db: Db;
const EMITTER_EVENT_NAME = "update";
constructor(db: Db, subscriber: Redis) {
export type SubscribeCallback = (tenant: Tenant) => void;
interface TenantUpdateMessage {
tenant: Tenant;
clientApplicationID: string;
}
// TenantCache provides an interface for retrieving tenant stored in local
// memory rather than grabbing it from the database every single call.
export default class TenantCache {
/**
* tenantsByID reference the tenants that have been cached/retrieved by ID.
*/
private tenantsByID = new DataLoader<string, Readonly<Tenant> | null>(ids => {
logger.debug({ ids: ids.length }, "now loading tenants");
return retrieveManyTenants(this.mongo, ids);
});
/**
* tenantsByDomain reference the tenants that have been cached/retrieved by
* Domain.
*/
private tenantsByDomain = new DataLoader<string, Readonly<Tenant> | null>(
domains => {
logger.debug({ domains: domains.length }, "now loading tenants");
return retrieveManyTenantsByDomain(this.mongo, domains);
}
);
/**
* Create a new client application ID. This prevents duplicated messages
* generated by this application from being handled as external messages
* as we should have already processed it.
*/
private clientApplicationID = uuid.v4();
private mongo: Db;
private emitter = new EventEmitter();
constructor(mongo: Db, subscriber: Redis) {
// Save the Db reference.
this.db = db;
// Prepare the list of all tenant's maintained by this instance.
this.tenants = new DataLoader(ids => retrieveManyTenants(db, ids));
// Subscribe to tenant notifications.
subscriber.subscribe(CacheUpdateChannel);
this.mongo = mongo;
// Attach to messages on this connection so we can receive updates when
// the tenant are changed.
subscriber.on("message", this.onMessage);
// Subscribe to tenant notifications.
subscriber.subscribe(TENANT_UPDATE_CHANNEL);
}
/**
@@ -37,13 +72,19 @@ export default class Cache {
*/
public async primeAll() {
// Grab all the tenants for this node.
const tenants = await retrieveAllTenants(this.db);
const tenants = await retrieveAllTenants(this.mongo);
// Clear out all the items in the cache.
this.tenants.clearAll();
this.tenantsByID.clearAll();
this.tenantsByDomain.clearAll();
// Prime the cache with each of these tenants.
tenants.forEach(tenant => this.tenants.prime(tenant.id, tenant));
tenants.forEach(tenant => {
this.tenantsByID.prime(tenant.id, tenant);
this.tenantsByDomain.prime(tenant.domain, tenant);
});
logger.debug({ tenants: tenants.length }, "primed tenants");
}
/**
@@ -54,26 +95,61 @@ export default class Cache {
message: string
): Promise<void> => {
// Only do things when the message is for tenant.
if (channel !== CacheUpdateChannel) {
if (channel !== TENANT_UPDATE_CHANNEL) {
return;
}
try {
// Updated tenant come from the messages.
const tenant: Tenant = JSON.parse(message);
const { tenant, clientApplicationID }: TenantUpdateMessage = JSON.parse(
message
);
// Check to see if this was the update issued by this instance.
if (clientApplicationID === this.clientApplicationID) {
// It was, so just return here, we already updated/handled it.
return;
}
logger.debug({ tenant_id: tenant.id }, "received updated tenant");
// Update the tenant cache.
this.tenants.clear(tenant.id).prime(tenant.id, tenant);
this.tenantsByID.clear(tenant.id).prime(tenant.id, tenant);
this.tenantsByDomain.clear(tenant.domain).prime(tenant.domain, tenant);
// Publish the event for the connected listeners.
this.emitter.emit(EMITTER_EVENT_NAME, tenant);
} catch (err) {
// FIXME: handle the error
logger.error(
{ err },
"an error occurred while trying to parse/prime the tenant/tenant cache"
);
}
};
public async retrieveByID(id: string): Promise<Readonly<Tenant> | null> {
return this.tenantsByID.load(id);
}
public async retrieveByDomain(
domain: string
): Promise<Readonly<Tenant> | null> {
return this.tenantsByDomain.load(domain);
}
/**
* retrieve returns a promise that will resolve to the tenant for Talk.
* This allows you to subscribe to new Tenant updates. This will also return
* a function that when called, unsubscribes you from updates.
*
* @param callback the function to be called when there is an updated Tenant.
*/
public async retrieve(id: string): Promise<Readonly<Tenant> | null> {
return this.tenants.load(id);
public subscribe(callback: SubscribeCallback) {
this.emitter.on(EMITTER_EVENT_NAME, callback);
// Return the unsubscribe function.
return () => {
this.emitter.removeListener(EMITTER_EVENT_NAME, callback);
};
}
/**
@@ -85,9 +161,23 @@ export default class Cache {
*/
public async update(conn: Redis, tenant: Tenant): Promise<void> {
// Update the tenant in the local cache.
this.tenants.clear(tenant.id).prime(tenant.id, tenant);
this.tenantsByID.clear(tenant.id).prime(tenant.id, tenant);
this.tenantsByDomain.clear(tenant.domain).prime(tenant.domain, tenant);
// Notify the other nodes about the tenant change.
await conn.publish(CacheUpdateChannel, JSON.stringify(tenant));
const message: TenantUpdateMessage = {
tenant,
clientApplicationID: this.clientApplicationID,
};
const subscribers = await conn.publish(
TENANT_UPDATE_CHANNEL,
JSON.stringify(message)
);
logger.debug({ tenant_id: tenant.id, subscribers }, "updated tenant");
// Publish the event for the connected listeners.
this.emitter.emit(EMITTER_EVENT_NAME, tenant);
}
}
+30
View File
@@ -0,0 +1,30 @@
import { Redis } from "ioredis";
import { Db } from "mongodb";
import {
Tenant,
updateTenant,
UpdateTenantInput,
} from "talk-server/models/tenant";
import TenantCache from "./cache";
export type UpdateTenant = UpdateTenantInput;
export async function update(
db: Db,
conn: Redis,
cache: TenantCache,
tenant: Tenant,
input: UpdateTenant
): Promise<Tenant | null> {
const updatedTenant = await updateTenant(db, tenant.id, input);
if (!updatedTenant) {
return null;
}
// Update the tenant cache.
await cache.update(conn, updatedTenant);
return updatedTenant;
}
+22
View File
@@ -0,0 +1,22 @@
import { get } from "lodash";
import { KarmaThresholds } from "talk-server/models/settings";
import { User } from "talk-server/models/user";
export const getCommentTrustScore = (user: User): number =>
get(user, "metadata.trust.comment.karma", 0);
export const isReliableCommenter = (
thresholds: KarmaThresholds,
user: User
): boolean | null => {
const score = getCommentTrustScore(user);
if (score >= thresholds.comment.reliable) {
return true;
} else if (score <= thresholds.comment.unreliable) {
return false;
}
return null;
};
+2
View File
@@ -2,8 +2,10 @@ import { Request } from "express";
import { Tenant } from "talk-server/models/tenant";
import { User } from "talk-server/models/user";
import TenantCache from "talk-server/services/tenant/cache";
export interface Request extends Request {
user?: User;
tenant?: Tenant;
tenantCache: TenantCache;
}