Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/reuse-quality.yml
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ jobs:
run: docker compose up -d --wait

- name: Start dev API
run: NODE_ENV=development NODE_CONFIG_DIR=api/config node api/index.ts &
run: NODE_ENV=development NODE_CONFIG_DIR=./api/config node api/index.ts &

- name: Wait for API to be ready
run: |
Expand Down
2 changes: 1 addition & 1 deletion api/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@
},
"dependencies": {
"@data-fair/lib-common-types": "^1.24.0",
"@data-fair/lib-express": "^1.22.5",
"@data-fair/lib-express": "^1.27.0",
"@data-fair/lib-node": "^2.12.1",
"@data-fair/lib-utils": "^1.14.0",
"@data-fair/lib-validation": "^1.0.2",
Expand Down
2 changes: 1 addition & 1 deletion api/src/events/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ router.get('', async (req, res, next) => {
query['resource.id'] = id
}

const project = mongoProjection(req.query.select, ['_search', 'htmlBody'])
const project = mongoProjection(req.query.select, ['_search', '_needsSearch', 'htmlBody'])

// implement a special pagination based on the fact that we always sort by date
const sort: Sort = { date: -1, _id: -1 }
Expand Down
61 changes: 61 additions & 0 deletions api/src/events/search-worker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,61 @@
// Rebuilds the search texts of the events flagged by the identity webhooks (_needsSearch),
// so that simple-directory gets its response without waiting for the text index.
// Runs in the api server, like the webhooks worker.

import config from '#config'
import Debug from 'debug'
import mongo from '#mongo'
import { internalError } from '@data-fair/lib-node/observer.js'
import locks from '@data-fair/lib-node/locks.js'
import { buildSearchTexts } from './operations.ts'

const debug = Debug('search-worker')
const batchSize = 1000
const searchProjection = { topic: 1, title: 1, body: 1, sender: 1, originator: 1, _needsSearch: 1 }

let loopPromise: Promise<void> | null = null
let stopped = false
let acquiredLock = false

const wait = () => new Promise(resolve => setTimeout(resolve, config.worker.loopInterval))

const loop = async () => {
// eslint-disable-next-line no-unmodified-loop-condition
while (!stopped) {
try {
if (!acquiredLock) {
acquiredLock = await locks.acquire('search-loop')
if (!acquiredLock) { await wait(); continue }
}
// only what buildSearchTexts reads: an htmlBody can be heavy
const events = await mongo.events.find({ _needsSearch: { $exists: true } }, { projection: searchProjection }).limit(batchSize).toArray()
if (!events.length) { await wait(); continue }
debug('rebuild the search texts of', events.length, 'events')
await mongo.events.bulkWrite(events.map(event => {
// a newer rename flagged the event again in the meantime: skipped, picked up by the next batch
const filter = { _id: event._id, _needsSearch: event._needsSearch }
try {
return { updateOne: { filter, update: { $set: { _search: buildSearchTexts(event, config.i18n.locales, config.i18n.defaultLocale) }, $unset: { _needsSearch: 1 } } } }
} catch (err) {
// like drainSearchIndex in data-fair: the event keeps its previous texts, but the flag is cleared
// or this event would come back first in every batch and block the others forever
internalError('search-loop-event', `failed to rebuild the search texts of event ${event._id} - ${(err as Error)?.stack || err}`)
return { updateOne: { filter, update: { $unset: { _needsSearch: 1 } } } }
}
}), { ordered: false })
} catch (err) {
internalError('search-loop', err)
await wait()
}
}
if (acquiredLock) await locks.release('search-loop')
}

export const start = () => {
loopPromise = loop()
}

export const stop = async () => {
stopped = true
await loopPromise
}
57 changes: 2 additions & 55 deletions api/src/identities/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,59 +3,6 @@

import config from '#config'
import { createIdentitiesRouter } from '@data-fair/lib-express/identities/index.js'
import mongo from '#mongo'
import { updateIdentity, deleteIdentity } from './service.ts'

export default createIdentitiesRouter(
config.secretKeys.identities,
// onUpdate
async (identity) => {
if (identity.type === 'user') {
await mongo.notifications.updateMany({ 'recipient.id': identity.id }, { $set: { 'recipient.name': identity.name } })
await mongo.subscriptions.updateMany({ 'recipient.id': identity.id }, { $set: { 'recipient.name': identity.name } })
}
await mongo.subscriptions.updateMany({ 'sender.type': identity.type, 'sender.id': identity.id }, { $set: { 'sender.name': identity.name } })
await mongo.pushSubscriptions.updateMany({ 'owner.type': identity.type, 'owner.id': identity.id }, { $set: { 'owner.name': identity.name } })
await mongo.webhookSubscriptions.updateMany({ 'sender.type': identity.type, 'sender.id': identity.id }, { $set: { 'sender.name': identity.name } })
await mongo.webhookSubscriptions.updateMany({ 'owner.type': identity.type, 'owner.id': identity.id }, { $set: { 'owner.name': identity.name } })
if (identity.departments) {
for (const department of identity.departments.filter(d => !!d.name)) {
await mongo.subscriptions.updateMany({ 'sender.type': identity.type, 'sender.id': identity.id, 'sender.department': department.id }, { $set: { 'sender.name': identity.name, 'sender.departmentName': department.name } })
await mongo.pushSubscriptions.updateMany({ 'owner.type': identity.type, 'owner.id': identity.id, 'owner.department': department.id }, { $set: { 'owner.name': identity.name, 'owner.departmentName': department.name } })
await mongo.webhookSubscriptions.updateMany({ 'sender.type': identity.type, 'sender.id': identity.id, 'sender.department': department.id }, { $set: { 'sender.name': identity.name, 'sender.departmentName': department.name } })
await mongo.webhookSubscriptions.updateMany({ 'owner.type': identity.type, 'owner.id': identity.id, 'owner.department': department.id }, { $set: { 'owner.name': identity.name, 'owner.departmentName': department.name } })
}
}

if (identity.type === 'user' && identity.organizations) {
const privateSubscriptionFilter = {
'recipient.id': identity.id,
visibility: { $ne: 'public' as const },
'sender.type': 'organization'
}
for await (const privateSubscription of mongo.subscriptions.find(privateSubscriptionFilter)) {
let userOrg = identity.organizations.find(o => o.id === privateSubscription.sender?.id && !o.department)
if (privateSubscription.sender?.department) {
userOrg = userOrg || identity.organizations.find(o => o.id === privateSubscription.sender?.id && o.department === privateSubscription.sender.department)
}
if (userOrg && privateSubscription.sender?.role && userOrg.role !== privateSubscription.sender.role && userOrg.role !== 'admin') {
userOrg = undefined
}
if (!userOrg) {
// console.log('remove private subscription that does not match user orgs anymore', identity, privateSubscription)
await mongo.subscriptions.deleteOne({ _id: privateSubscription._id })
}
}
}
},
// onDelete
async (identity) => {
if (identity.type === 'user') {
await mongo.notifications.deleteMany({ 'recipient.id': identity.id })
await mongo.subscriptions.deleteMany({ 'recipient.id': identity.id })
}
await mongo.subscriptions.deleteMany({ 'sender.type': identity.type, 'sender.id': identity.id })
await mongo.pushSubscriptions.deleteMany({ 'owner.type': identity.type, 'owner.id': identity.id })
await mongo.webhookSubscriptions.deleteMany({ 'owner.type': identity.type, 'owner.id': identity.id })
await mongo.webhookSubscriptions.deleteMany({ 'sender.type': identity.type, 'sender.id': identity.id })
}
)
export default createIdentitiesRouter(config.secretKeys.identities, updateIdentity, deleteIdentity)
103 changes: 103 additions & 0 deletions api/src/identities/service.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
// Synchronize the copies of identity data (names on senders, recipients, owners, originators)
// with the users/organizations directory, and remove them when an identity is deleted.
// Everything is done with bulk updates: an organization can own tens of thousands of events and
// simple-directory waits for the response. Names are also part of the events search texts: the
// events are only flagged here (_needsSearch) and the search worker rebuilds the texts.

import { randomUUID } from 'node:crypto'
import type { IdentityUpdate, IdentityDelete } from '@data-fair/lib-express/identities/index.js'
import mongo from '#mongo'

export const updateIdentity = async (identity: IdentityUpdate) => {
const { type, id, name, departments } = identity

if (type === 'user') {
await mongo.notifications.updateMany({ 'recipient.id': id }, { $set: { 'recipient.name': name } })
await mongo.subscriptions.updateMany({ 'recipient.id': id }, { $set: { 'recipient.name': name } })
}
// notifications are snapshots taken at delivery, their sender name is not rewritten (no index on sender)
await mongo.subscriptions.updateMany({ 'sender.type': type, 'sender.id': id }, { $set: { 'sender.name': name } })
await mongo.pushSubscriptions.updateMany({ 'owner.type': type, 'owner.id': id }, { $set: { 'owner.name': name } })
await mongo.webhookSubscriptions.updateMany({ 'sender.type': type, 'sender.id': id }, { $set: { 'sender.name': name } })
await mongo.webhookSubscriptions.updateMany({ 'owner.type': type, 'owner.id': id }, { $set: { 'owner.name': name } })
if (departments) {
for (const department of departments.filter(d => !!d.name)) {
await mongo.subscriptions.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.department': department.id }, { $set: { 'sender.name': name, 'sender.departmentName': department.name } })
await mongo.pushSubscriptions.updateMany({ 'owner.type': type, 'owner.id': id, 'owner.department': department.id }, { $set: { 'owner.name': name, 'owner.departmentName': department.name } })
await mongo.webhookSubscriptions.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.department': department.id }, { $set: { 'sender.name': name, 'sender.departmentName': department.name } })
await mongo.webhookSubscriptions.updateMany({ 'owner.type': type, 'owner.id': id, 'owner.department': department.id }, { $set: { 'owner.name': name, 'owner.departmentName': department.name } })
}
// the directory sends the complete list of departments: a department missing from it was
// deleted, what it owns keeps the id (still reachable by the organization admins) but not the name
const deletedDepartment = { $exists: true, $nin: departments.map(d => d.id) }
await mongo.subscriptions.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.department': deletedDepartment }, { $unset: { 'sender.departmentName': 1 } })
await mongo.pushSubscriptions.updateMany({ 'owner.type': type, 'owner.id': id, 'owner.department': deletedDepartment }, { $unset: { 'owner.departmentName': 1 } })
await mongo.webhookSubscriptions.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.department': deletedDepartment }, { $unset: { 'sender.departmentName': 1 } })
await mongo.webhookSubscriptions.updateMany({ 'owner.type': type, 'owner.id': id, 'owner.department': deletedDepartment }, { $unset: { 'owner.departmentName': 1 } })
}

// events: as sender, and as the user or organization that triggered them
// only the events whose name actually changes: simple-directory also posts on every membership change.
// A fresh _needsSearch per update: the search worker only clears the value it read.
await mongo.events.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.name': { $ne: name } }, { $set: { 'sender.name': name, _needsSearch: randomUUID() } })
if (departments) {
for (const department of departments.filter(d => !!d.name)) {
await mongo.events.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.department': department.id, 'sender.departmentName': { $ne: department.name } }, { $set: { 'sender.departmentName': department.name } })
await mongo.events.updateMany({ 'originator.organization.id': id, 'originator.organization.department': department.id, 'originator.organization.departmentName': { $ne: department.name } }, { $set: { 'originator.organization.departmentName': department.name } })
}
const deletedDepartment = { $exists: true, $nin: departments.map(d => d.id) }
await mongo.events.updateMany({ 'sender.type': type, 'sender.id': id, 'sender.department': deletedDepartment }, { $unset: { 'sender.departmentName': 1 } })
await mongo.events.updateMany({ 'originator.organization.id': id, 'originator.organization.department': deletedDepartment }, { $unset: { 'originator.organization.departmentName': 1 } })
}
if (type === 'user') {
await mongo.events.updateMany({ 'originator.user.id': id, 'originator.user.name': { $ne: name } }, { $set: { 'originator.user.name': name, _needsSearch: randomUUID() } })
} else {
await mongo.events.updateMany({ 'originator.organization.id': id, 'originator.organization.name': { $ne: name } }, { $set: { 'originator.organization.name': name, _needsSearch: randomUUID() } })
}

if (type === 'user' && identity.organizations) {
const privateSubscriptionFilter = {
'recipient.id': id,
visibility: { $ne: 'public' as const },
'sender.type': 'organization'
}
for await (const privateSubscription of mongo.subscriptions.find(privateSubscriptionFilter)) {
let userOrg = identity.organizations.find(o => o.id === privateSubscription.sender?.id && !o.department)
if (privateSubscription.sender?.department) {
userOrg = userOrg || identity.organizations.find(o => o.id === privateSubscription.sender?.id && o.department === privateSubscription.sender.department)
}
if (userOrg && privateSubscription.sender?.role && userOrg.role !== privateSubscription.sender.role && userOrg.role !== 'admin') {
userOrg = undefined
}
if (!userOrg) {
// remove private subscription that does not match user orgs anymore
await mongo.subscriptions.deleteOne({ _id: privateSubscription._id })
}
}
}
}

export const deleteIdentity = async (identity: IdentityDelete) => {
const { type, id } = identity

if (type === 'user') {
await mongo.notifications.deleteMany({ 'recipient.id': id })
await mongo.subscriptions.deleteMany({ 'recipient.id': id })
await mongo.pointers.deleteMany({ 'recipient.id': id })
}
await mongo.subscriptions.deleteMany({ 'sender.type': type, 'sender.id': id })
await mongo.pushSubscriptions.deleteMany({ 'owner.type': type, 'owner.id': id })
await mongo.webhookSubscriptions.deleteMany({ 'owner.type': type, 'owner.id': id })
await mongo.webhookSubscriptions.deleteMany({ 'sender.type': type, 'sender.id': id })
// pending or failed webhooks of the deleted subscriptions
await mongo.webhooks.deleteMany({ 'owner.type': type, 'owner.id': id })
await mongo.webhooks.deleteMany({ 'sender.type': type, 'sender.id': id })

// the events of the identity are its own feed, nobody else can read them
await mongo.events.deleteMany({ 'sender.type': type, 'sender.id': id })
// the events a user triggered on other feeds keep the trace of the action without the person:
// only the id remains (pseudonymized), an organization is not personal data and is left as is
if (type === 'user') {
await mongo.events.updateMany({ 'originator.user.id': id }, { $set: { _needsSearch: randomUUID() }, $unset: { 'originator.user.name': 1, 'originator.user.email': 1 } })
}
}
9 changes: 8 additions & 1 deletion api/src/mongo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,12 @@ export class EventsMongo {
'main-keys': [
{ 'sender.type': 1, 'sender.id': 1, '_search.text': 'text', date: -1 },
{ default_language: config.i18n.defaultLocale }
]
],
// identity webhooks rewrite the events triggered by a user or an organization
'originator-user': [{ 'originator.user.id': 1 }, { sparse: true }],
'originator-organization': [{ 'originator.organization.id': 1 }, { sparse: true }],
// the events waiting for the search worker, empty most of the time
'needs-search': [{ _needsSearch: 1 }, { partialFilterExpression: { _needsSearch: { $exists: true } } }]
},
subscriptions: {
'main-keys': [
Expand All @@ -75,6 +80,8 @@ export class EventsMongo {
},
webhooks: {
'main-keys': { 'owner.type': 1, 'owner.id': 1, 'subscription._id': 1, 'notification.date': 1 },
// identity webhooks drop the webhooks of a deleted sender
'sender-keys': { 'sender.type': 1, 'sender.id': 1 },
'loop-keys': { status: 1, nextAttempt: 1 },
'coalesce-keys': { 'subscription._id': 1, 'notification.topic.key': 1 }
},
Expand Down
3 changes: 3 additions & 0 deletions api/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import { createHttpTerminator } from 'http-terminator'
import app from './app.ts'
import config from '#config'
import * as webhooksWorker from './webhooks/worker.ts'
import * as searchWorker from './events/search-worker.ts'
import * as pushService from './push/service.ts'

const server = createServer(app)
Expand Down Expand Up @@ -41,6 +42,7 @@ export const start = async () => {
await wsEmitter.init(mongo.db)
await pushService.init()
await webhooksWorker.start()
searchWorker.start()

server.listen(config.port)
await new Promise(resolve => server.once('listening', resolve))
Expand All @@ -51,6 +53,7 @@ export const start = async () => {
export const stop = async () => {
await httpTerminator.terminate()
await webhooksWorker.stop()
await searchWorker.stop()
await wsServer.stop()
if (config.observer.active) await stopObserver()
await locks.stop()
Expand Down
3 changes: 2 additions & 1 deletion api/types/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,4 +10,5 @@ export type { DeviceRegistration } from './device-registration/index.js'

export type FullEvent = Event & Required<Pick<Event, 'visibility'>>
export type LocalizedEvent = Omit<FullEvent, 'title' | 'body' | 'htmlBody'> & { title: string, body?: string, htmlBody?: string }
export type SearchableEvent = FullEvent & { _search: { language: string, text: string }[] }
// _needsSearch: set by the identity webhooks when a name changed, cleared by the search worker
export type SearchableEvent = FullEvent & { _search: { language: string, text: string }[], _needsSearch?: string }
2 changes: 2 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,8 @@ services:
DEBUG: session
ADMINS: '["superadmin@test.com"]'
PUBLIC_URL: http://${DEV_HOST}:${NGINX_PORT}/simple-directory
# in dev simple-directory only notifies this service, directly (the shared identities router refuses calls through the proxy)
IDENTITIES_WEBHOOKS: '[{"base":"http://localhost:${DEV_API_PORT}/api/identities","key":"SECRET_IDENTITIES"}]'
MAILDEV_ACTIVE: "true"
STORAGE_TYPE: file
ROLES_DEFAULTS: '["admin", "contrib", "user"]'
Expand Down
6 changes: 3 additions & 3 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

3 changes: 2 additions & 1 deletion package.json
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,7 @@
},
"relativeDependencies": {
"@data-fair/lib-express": "../lib/packages/express",
"@data-fair/lib-vue": "../lib/packages/vue"
"@data-fair/lib-vue": "../lib/packages/vue",
"@data-fair/lib-vuetify": "../lib/packages/vuetify"
}
}
Loading
Loading