auth/src/sync.js
2026-07-02 00:23:59 +00:00

104 lines
4 KiB
JavaScript

import { pool } from './db.js'
import { findEmployeeById, findEmployeeByEmail, getEmployeeDepartments, getSyncConfig } from './workforce.js'
export async function syncAllWorkforceUsers() {
const { rows: users } = await pool.query(
`SELECT id, email, name, workforce_user_id FROM users WHERE workforce_user_id IS NOT NULL AND active = true`
)
let deactivated = 0, emailUpdated = 0, rolesUpdated = 0
for (const user of users) {
try {
// Prefer ID lookup — survives email changes in Workforce.
// Fall back to email lookup if the API doesn't support GET /users/:id.
let wfUser = await findEmployeeById(user.workforce_user_id)
if (!wfUser) wfUser = await findEmployeeByEmail(user.email)
if (!wfUser || String(wfUser.id) !== user.workforce_user_id) {
// Not found in Workforce — deactivate
await pool.query('UPDATE users SET active = false WHERE id = $1', [user.id])
deactivated++
continue
}
// Sync email if changed
const wfEmail = wfUser.email?.toLowerCase().trim()
if (wfEmail && wfEmail !== user.email) {
await pool.query('UPDATE users SET email = $1 WHERE id = $2', [wfEmail, user.id])
emailUpdated++
}
// Sync department-mapped roles
const depts = await getEmployeeDepartments(wfUser.id)
const deptIds = depts.map(d => String(d.id))
// Get current dept-mapped role_ids for this user
const { rows: currentRoles } = await pool.query(
`SELECT ur.role_id, wdr.department_id
FROM user_roles ur
JOIN workforce_department_roles wdr ON wdr.role_id = ur.role_id
WHERE ur.user_id = $1 AND ur.source = 'workforce_department'`,
[user.id]
)
const currentDeptIds = currentRoles.map(r => r.department_id)
const added = deptIds.filter(id => !currentDeptIds.includes(id))
const removed = currentDeptIds.filter(id => !deptIds.includes(id))
if (added.length || removed.length) {
// Remove roles for departments the user is no longer in
if (removed.length) {
await pool.query(
`DELETE FROM user_roles WHERE user_id = $1 AND source = 'workforce_department'
AND role_id IN (
SELECT role_id FROM workforce_department_roles WHERE department_id = ANY($2)
)`,
[user.id, removed]
)
}
// Add roles for new departments
for (const deptId of added) {
await pool.query(
`INSERT INTO user_roles (user_id, role_id, source)
SELECT $1, role_id, 'workforce_department'
FROM workforce_department_roles WHERE department_id = $2
ON CONFLICT DO NOTHING`,
[user.id, deptId]
)
}
rolesUpdated++
}
} catch (err) {
console.error(`Workforce sync error for user ${user.id} (${user.email}):`, err.message)
}
}
// Clean up expired pending registrations older than 24h
await pool.query(`DELETE FROM pending_registrations WHERE expires_at < NOW() - INTERVAL '24 hours'`)
const summary = { checked: users.length, deactivated, emailUpdated, rolesUpdated }
console.log('Workforce sync complete:', summary)
return summary
}
export function startSyncJob() {
getSyncConfig()
.then(({ syncHours }) => {
if (!syncHours || syncHours <= 0) {
console.log('Workforce sync disabled (sync_hours = 0)')
return
}
const intervalMs = syncHours * 60 * 60_000
// First run 60s after startup to let DB settle
setTimeout(() => {
syncAllWorkforceUsers().catch(err => console.error('Workforce sync failed:', err.message))
setInterval(
() => syncAllWorkforceUsers().catch(err => console.error('Workforce sync failed:', err.message)),
intervalMs
)
}, 60_000)
console.log(`Workforce sync scheduled every ${syncHours}h`)
})
.catch(err => console.warn('Could not read Workforce sync config (settings not configured?):', err.message))
}