459 lines
21 KiB
Python
459 lines
21 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""Sync queue: user changes land in this table and are pushed to OdooshCN asynchronously.
|
|
|
|
Why the hooks do not call the API directly: that would make user management here depend
|
|
on the platform being up. Writing a queue row is a local insert and never fails.
|
|
|
|
Three paths, each a fallback for the previous one:
|
|
* kick_async - a dedicated thread right after the transaction commits, so the user
|
|
barely notices any delay.
|
|
* cron_tick - every minute, in case the thread was killed with the worker. It is also
|
|
the only thing that drives the retry back-off: nothing else comes back
|
|
for a row whose next_retry_at is four hours away.
|
|
* cron_scan - folded into that same tick, at most every SCAN_INTERVAL_DEFAULT minutes.
|
|
Walks users by write_date and compares a payload fingerprint, catching
|
|
anything the hooks missed (changes made outside the ORM, or made while
|
|
the module was disabled).
|
|
|
|
The scan shares the queue cron's slot instead of owning one. Both jobs are cheap in
|
|
themselves, but Odoo's cron pool is shared by the whole database (max_cron_threads, 2 by
|
|
default), so a slot taken every few minutes costs more than the work does.
|
|
|
|
Retry policy: temporary failures back off 1, 5, 15, 60 and 240 minutes, then the row is
|
|
marked failed and waits for a human. Permanent failures are not retried at all.
|
|
"""
|
|
|
|
import json
|
|
import logging
|
|
import threading
|
|
from datetime import timedelta
|
|
|
|
from odoo import SUPERUSER_ID, _, api, fields, models
|
|
from odoo.exceptions import UserError
|
|
from odoo.modules.registry import Registry
|
|
|
|
from .odoosh_client import OdooshPermanent, OdooshRetryable
|
|
|
|
_logger = logging.getLogger(__name__)
|
|
|
|
# Minutes to wait after the n-th failure; once exhausted the row is marked failed
|
|
BACKOFF_MINUTES = [1, 5, 15, 60, 240]
|
|
BATCH_SIZE = 50 # rows processed per queue round
|
|
SCAN_BATCH = 200 # users inspected per scan round
|
|
SCAN_INTERVAL_DEFAULT = 30 # minutes between two scans; the parameter set to 0 disables it
|
|
PARAM_LAST_SCAN = 'odoosh.last_scan_at' # cursor: user write_date reached by the scan
|
|
PARAM_SCAN_RUN = 'odoosh.last_scan_run' # when the scan last ran, to space the runs out
|
|
PARAM_SCAN_INTERVAL = 'odoosh.scan_interval_minutes'
|
|
|
|
# At most one push thread per database: a second request finding the lock taken just
|
|
# returns, because the running thread re-queries the queue until it is empty.
|
|
_workers = {}
|
|
_workers_guard = threading.Lock()
|
|
|
|
|
|
class OdooshSyncLog(models.Model):
|
|
_name = 'odoosh.sync.log'
|
|
_description = 'OdooshCN Sync Queue'
|
|
_order = 'id desc'
|
|
_rec_name = 'display_name'
|
|
|
|
# ondelete must be `set null`: when the Odoo user is deleted this row still has to
|
|
# travel to the platform to delete the account there.
|
|
user_id = fields.Many2one('res.users', string='User', ondelete='set null', index=True)
|
|
login = fields.Char(string='Login', help="Kept so the row still identifies someone after the user is deleted")
|
|
external_id = fields.Char(string='External ID', required=True, index=True,
|
|
help="The stable key sent to OdooshCN, i.e. the user ID in this system")
|
|
operation = fields.Selection([
|
|
('upsert', 'Create / Update'),
|
|
('deactivate', 'Deactivate'),
|
|
('delete', 'Delete'),
|
|
], string='Operation', required=True, default='upsert')
|
|
state = fields.Selection([
|
|
('pending', 'Pending'),
|
|
('done', 'Done'),
|
|
('skipped', 'Skipped (no change)'),
|
|
('failed', 'Failed'),
|
|
('cancelled', 'Cancelled'),
|
|
], string='Status', default='pending', required=True, index=True)
|
|
force = fields.Boolean(string='Force push',
|
|
help="Ignore the payload fingerprint and push even when nothing changed")
|
|
attempts = fields.Integer(string='Attempts', default=0)
|
|
next_retry_at = fields.Datetime(string='Next retry', index=True)
|
|
last_error = fields.Text(string='Last error')
|
|
error_code = fields.Char(string='Error code',
|
|
help="The platform's stable error code, used to decide retry versus manual handling")
|
|
payload_preview = fields.Text(string='Payload', help="What was sent last time; contains no secrets")
|
|
result = fields.Text(string='Response')
|
|
synced_at = fields.Datetime(string='Completed on')
|
|
display_name = fields.Char(compute='_compute_display_name')
|
|
|
|
@api.depends('login', 'operation', 'state')
|
|
def _compute_display_name(self):
|
|
ops = dict(self._fields['operation'].selection)
|
|
for rec in self:
|
|
rec.display_name = '%s - %s' % (rec.login or _('Unknown user'), ops.get(rec.operation, rec.operation))
|
|
|
|
# ------------------------------------------------------------------ enqueue
|
|
@api.model
|
|
def enqueue(self, users, operation='upsert', force=False):
|
|
"""Queue a batch of users. An existing pending row for the same user and operation
|
|
is reused (back-off reset, latest data resent) so renaming twice does not pile up
|
|
rows. One query for the whole batch, never one per user."""
|
|
users = users.sudo()
|
|
if not users:
|
|
return self.browse()
|
|
pending = self.sudo().search([
|
|
('user_id', 'in', users.ids),
|
|
('operation', '=', operation),
|
|
('state', '=', 'pending'),
|
|
])
|
|
by_user = {p.user_id.id: p for p in pending}
|
|
if pending:
|
|
vals = {'next_retry_at': False, 'attempts': 0}
|
|
if force:
|
|
vals['force'] = True
|
|
pending.write(vals)
|
|
new_vals = [{
|
|
'user_id': u.id,
|
|
'login': u.login,
|
|
'external_id': str(u.id),
|
|
'operation': operation,
|
|
'state': 'pending',
|
|
'force': force,
|
|
} for u in users if u.id not in by_user]
|
|
created = self.sudo().create(new_vals) if new_vals else self.browse()
|
|
return pending | created
|
|
|
|
@api.model
|
|
def enqueue_delete(self, users):
|
|
"""Deleting an Odoo user: record the identity now, because after `unlink` the
|
|
user record is gone and the row has to survive on its own."""
|
|
users = users.sudo()
|
|
if not users:
|
|
return self.browse()
|
|
return self.sudo().create([{
|
|
'user_id': u.id,
|
|
'login': u.login,
|
|
'external_id': str(u.id),
|
|
'operation': 'delete',
|
|
'state': 'pending',
|
|
} for u in users])
|
|
|
|
# ------------------------------------------------------------------ background thread
|
|
@api.model
|
|
def kick_async(self, dbname=None, scan_first=False):
|
|
"""Called after commit: drain the queue in a dedicated thread. One thread per
|
|
database. `scan_first` runs a user scan first, used after the settings change.
|
|
|
|
The thread owns its cursor and environment, fully isolated from the request that
|
|
started it; anything it raises is logged and cannot affect committed data."""
|
|
dbname = dbname or self.env.cr.dbname
|
|
with _workers_guard:
|
|
lock = _workers.setdefault(dbname, threading.Lock())
|
|
if not lock.acquire(blocking=False):
|
|
return False # a thread is already running and will pick up the new rows
|
|
|
|
def run():
|
|
threading.current_thread().dbname = dbname # makes the log lines carry the db name
|
|
try:
|
|
registry = Registry(dbname)
|
|
if scan_first:
|
|
with registry.cursor() as cr:
|
|
api.Environment(cr, SUPERUSER_ID, {})['odoosh.sync.log'].cron_scan(kick=False)
|
|
rounds = 0
|
|
while rounds < 20: # guard: hand back to the cron if rows keep pouring in
|
|
rounds += 1
|
|
with registry.cursor() as cr:
|
|
env = api.Environment(cr, SUPERUSER_ID, {})
|
|
n = env['odoosh.sync.log'].cron_process()
|
|
if n == 0:
|
|
break
|
|
except Exception: # noqa: BLE001
|
|
_logger.exception("OdooshCN push thread died (the cron will take over)")
|
|
finally:
|
|
lock.release()
|
|
|
|
threading.Thread(target=run, name='odoosh-sync-%s' % dbname, daemon=True).start()
|
|
return True
|
|
|
|
# ------------------------------------------------------------------ the cron
|
|
@api.model
|
|
def cron_tick(self):
|
|
"""The module's only frequent cron: drain the queue, and fold the user scan into the
|
|
same run every `odoosh.scan_interval_minutes` minutes (0 turns the scan off).
|
|
|
|
Scanning first and pushing afterwards means whatever the scan finds leaves in this
|
|
tick instead of waiting for the next one."""
|
|
if not self.env['odoosh.client']._enabled():
|
|
return 0
|
|
if self._scan_due():
|
|
self.cron_scan(kick=False) # kick_async would only duplicate the push below
|
|
return self.cron_process()
|
|
|
|
@api.model
|
|
def _scan_due(self):
|
|
"""Whether this tick also scans. The timestamp is written before the scan runs, so a
|
|
scan that crashes waits for the next interval instead of retrying every minute."""
|
|
icp = self.env['ir.config_parameter'].sudo()
|
|
try:
|
|
minutes = int(icp.get_param(PARAM_SCAN_INTERVAL) or SCAN_INTERVAL_DEFAULT)
|
|
except (TypeError, ValueError):
|
|
minutes = SCAN_INTERVAL_DEFAULT
|
|
if minutes <= 0:
|
|
return False
|
|
now = fields.Datetime.now()
|
|
last = icp.get_param(PARAM_SCAN_RUN)
|
|
if last:
|
|
try:
|
|
if now - fields.Datetime.to_datetime(last) < timedelta(minutes=minutes):
|
|
return False
|
|
except ValueError:
|
|
pass # unreadable value: scan now and write a good one back
|
|
icp.set_param(PARAM_SCAN_RUN, fields.Datetime.to_string(now))
|
|
return True
|
|
|
|
# ------------------------------------------------------------------ queue processing
|
|
@api.model
|
|
def cron_process(self, limit=BATCH_SIZE):
|
|
"""Process the rows that are due; returns how many were handled, failures included."""
|
|
client = self.env['odoosh.client']
|
|
if not client._enabled():
|
|
return 0
|
|
now = fields.Datetime.now()
|
|
records = self.sudo().search([
|
|
('state', '=', 'pending'),
|
|
'|', ('next_retry_at', '=', False), ('next_retry_at', '<=', now),
|
|
], order='id asc', limit=limit)
|
|
if not records:
|
|
return 0
|
|
ctx = self.env['res.users']._odoosh_scope_ctx() # shared by the whole batch
|
|
for rec in records:
|
|
rec._process_one(ctx)
|
|
return len(records)
|
|
|
|
def _process_one(self, ctx=None):
|
|
"""Process one row; returns whether it succeeded (skipped counts as success).
|
|
Each row commits on its own so one failure never rolls back the others."""
|
|
self.ensure_one()
|
|
client = self.env['odoosh.client']
|
|
user = self.user_id.sudo() if self.user_id else self.env['res.users']
|
|
try:
|
|
if self.operation == 'delete':
|
|
payload = {'external_id': self.external_id}
|
|
result = client.delete_user(self.external_id)
|
|
self._finish('done', payload, result)
|
|
if user.exists():
|
|
user._odoosh_mark('none')
|
|
return True
|
|
|
|
if self.operation == 'deactivate':
|
|
payload = {'external_id': self.external_id, 'stop_envs': True}
|
|
result = client.deactivate_user(self.external_id, stop_envs=True)
|
|
self._finish('done', payload, result)
|
|
if user.exists():
|
|
user._odoosh_mark('none')
|
|
return True
|
|
|
|
if not user.exists():
|
|
self._mark_failed(_("The user was deleted; there is nothing left to sync"), 'user_gone')
|
|
return False
|
|
payload = user._odoosh_payload(ctx)
|
|
digest = user._odoosh_payload_hash(payload)
|
|
if not self.force and user.odoosh_sync_state == 'synced' and user.odoosh_payload_hash == digest:
|
|
self._finish('skipped', payload, {'note': 'unchanged'}) # nothing changed, leave the platform alone
|
|
return True
|
|
try:
|
|
result = client.upsert_user(payload)
|
|
except OdooshPermanent as err:
|
|
# Platform username already taken (typically Odoo's `admin` colliding with the
|
|
# platform's own admin): retry once with a "-<odoo user id>" suffix and remember it.
|
|
if err.code != 'username_taken' or (user.odoosh_username or '').strip():
|
|
raise
|
|
alt = ('%s-%s' % (payload['username'], user.id))[:64]
|
|
payload['username'] = alt
|
|
if payload.get('provision'):
|
|
payload['provision']['name'] = alt
|
|
result = client.upsert_user(payload)
|
|
user.with_context(odoosh_no_sync=True).write({'odoosh_username': alt})
|
|
digest = user._odoosh_payload_hash(payload)
|
|
self._finish('done', payload, result)
|
|
user._odoosh_mark('synced', payload_hash=digest, result=result)
|
|
return True
|
|
except OdooshPermanent as err:
|
|
self._mark_failed(err.message, err.code)
|
|
if user.exists():
|
|
user._odoosh_mark('failed', error='[%s] %s' % (err.code, err.message))
|
|
return False
|
|
except OdooshRetryable as err:
|
|
self._mark_retry(err.message, err.code)
|
|
if user.exists():
|
|
user._odoosh_mark('pending', error='[%s] %s' % (err.code, err.message))
|
|
return False
|
|
except Exception as err: # noqa: BLE001 - unexpected errors are retried too, never fatal
|
|
_logger.exception("OdooshCN sync crashed on log#%s", self.id)
|
|
self._mark_retry(str(err), 'unexpected')
|
|
return False
|
|
finally:
|
|
# Commit per row: progress survives the thread or worker being killed mid-way
|
|
self.env.cr.commit()
|
|
|
|
def _finish(self, state, payload, result):
|
|
self.write({
|
|
'state': state,
|
|
'attempts': self.attempts + 1,
|
|
'synced_at': fields.Datetime.now(),
|
|
'payload_preview': self._pretty(payload),
|
|
'result': self._pretty(result),
|
|
'last_error': False,
|
|
'error_code': False,
|
|
'next_retry_at': False,
|
|
'force': False,
|
|
})
|
|
|
|
def _mark_retry(self, message, code):
|
|
attempts = self.attempts + 1
|
|
if attempts > len(BACKOFF_MINUTES):
|
|
return self._mark_failed(
|
|
_("Gave up after %(n)s attempts. Last error: %(err)s", n=attempts - 1, err=message), code)
|
|
delay = BACKOFF_MINUTES[attempts - 1]
|
|
self.write({
|
|
'attempts': attempts,
|
|
'last_error': message,
|
|
'error_code': code,
|
|
'next_retry_at': fields.Datetime.now() + timedelta(minutes=delay),
|
|
})
|
|
_logger.warning("OdooshCN sync postponed log#%s (%s), retry in %s min: %s", self.id, code, delay, message)
|
|
|
|
def _mark_failed(self, message, code):
|
|
self.write({
|
|
'state': 'failed',
|
|
'attempts': self.attempts + 1,
|
|
'last_error': message,
|
|
'error_code': code,
|
|
'next_retry_at': False,
|
|
})
|
|
_logger.error("OdooshCN sync failed log#%s (%s): %s", self.id, code, message)
|
|
|
|
@staticmethod
|
|
def _pretty(data):
|
|
try:
|
|
return json.dumps(data, ensure_ascii=False, indent=2)
|
|
except (TypeError, ValueError):
|
|
return str(data)
|
|
|
|
# ------------------------------------------------------------------ periodic scan
|
|
@api.model
|
|
def cron_scan(self, limit=SCAN_BATCH, kick=True):
|
|
"""Walk users by write_date to catch whatever the hooks missed. Only rows changed
|
|
since the last scan are read, never the whole table.
|
|
|
|
Cost control: one indexed query, one in-memory fingerprint comparison per user,
|
|
at most SCAN_BATCH users per round; the cursor only advances to the last user
|
|
actually inspected, so nobody is skipped."""
|
|
if not self.env['odoosh.client']._enabled():
|
|
return 0
|
|
icp = self.env['ir.config_parameter'].sudo()
|
|
last = icp.get_param(PARAM_LAST_SCAN)
|
|
Users = self.env['res.users'].sudo()
|
|
ctx = Users._odoosh_scope_ctx()
|
|
|
|
domain = [('odoosh_sync_enabled', '=', True), ('share', '=', False)]
|
|
if last:
|
|
domain.append(('write_date', '>', last))
|
|
# Archived users matter too, so active_test is off; deactivation is only sent for
|
|
# accounts the platform already knows.
|
|
users = Users.with_context(active_test=False).search(domain, order='write_date asc, id asc', limit=limit)
|
|
if not users:
|
|
icp.set_param(PARAM_LAST_SCAN, fields.Datetime.to_string(fields.Datetime.now()))
|
|
return 0
|
|
|
|
to_upsert = Users.browse()
|
|
to_deactivate = Users.browse()
|
|
in_scope = users._odoosh_filter_in_scope(ctx)
|
|
for u in users:
|
|
if u not in in_scope:
|
|
continue # switch off, portal or built-in account
|
|
if not u.active:
|
|
if u.odoosh_sync_state in ('synced', 'pending', 'failed'):
|
|
to_deactivate |= u
|
|
continue
|
|
if u.odoosh_sync_state != 'synced' or u.odoosh_payload_hash != u._odoosh_payload_hash(ctx=ctx):
|
|
to_upsert |= u
|
|
|
|
if to_upsert:
|
|
self.enqueue(to_upsert, 'upsert')
|
|
to_upsert._odoosh_mark('pending')
|
|
if to_deactivate:
|
|
self.enqueue(to_deactivate, 'deactivate')
|
|
|
|
# Cursor: when the batch was full, only advance to its last write_date
|
|
cursor = users[-1].write_date if len(users) >= limit else fields.Datetime.now()
|
|
icp.set_param(PARAM_LAST_SCAN, fields.Datetime.to_string(cursor))
|
|
n = len(to_upsert) + len(to_deactivate)
|
|
if n:
|
|
_logger.info("OdooshCN scan: inspected %s users, queued %s changes", len(users), n)
|
|
self.env.cr.commit()
|
|
if kick:
|
|
self.kick_async()
|
|
return n
|
|
|
|
# ------------------------------------------------------------------ manual actions
|
|
def action_retry(self):
|
|
"""Put failed or cancelled rows back in the queue and process them at once."""
|
|
self.write({'state': 'pending', 'attempts': 0, 'next_retry_at': False, 'force': True})
|
|
ctx = self.env['res.users']._odoosh_scope_ctx()
|
|
for rec in self:
|
|
rec._process_one(ctx)
|
|
return True
|
|
|
|
def action_cancel(self):
|
|
self.filtered(lambda r: r.state in ('pending', 'failed')).write({'state': 'cancelled', 'next_retry_at': False})
|
|
return True
|
|
|
|
# ------------------------------------------------------------------ reconciliation
|
|
@api.model
|
|
def cron_reconcile(self):
|
|
"""Daily reconciliation: send the full list of users who should have an account,
|
|
so the platform can deactivate the ones that are gone.
|
|
|
|
An incomplete list would deactivate everyone, so an empty list is skipped here and
|
|
the platform applies a ratio guard of its own."""
|
|
client = self.env['odoosh.client']
|
|
if not client._enabled():
|
|
return False
|
|
users = self.env['res.users'].sudo().search([('active', '=', True)])._odoosh_filter_in_scope()
|
|
if not users:
|
|
_logger.warning("OdooshCN reconciliation skipped: the local list is empty")
|
|
return False
|
|
try:
|
|
result = client.reconcile([u.id for u in users], stop_envs=True, confirm=False)
|
|
_logger.info("OdooshCN reconciliation done: %s users sent, %s deactivated on the platform",
|
|
len(users), result.get('deactivated'))
|
|
except OdooshPermanent as err:
|
|
# The ratio guard lands here: never auto-confirm, let an administrator decide
|
|
_logger.warning("OdooshCN reconciliation refused (%s): %s", err.code, err.message)
|
|
except OdooshRetryable as err:
|
|
_logger.warning("OdooshCN reconciliation postponed (%s): %s", err.code, err.message)
|
|
return True
|
|
|
|
@api.model
|
|
def action_sync_all(self):
|
|
"""Re-queue every user in scope, ignoring fingerprints."""
|
|
users = self.env['res.users'].sudo().search([])._odoosh_filter_in_scope()
|
|
if not users:
|
|
raise UserError(_("No user is in sync scope. Check the settings and the per-user switch."))
|
|
self.enqueue(users, 'upsert', force=True)
|
|
users._odoosh_mark('pending')
|
|
self.env.cr.commit()
|
|
self.kick_async()
|
|
return {
|
|
'type': 'ir.actions.client',
|
|
'tag': 'display_notification',
|
|
'params': {
|
|
'title': _("Queued"),
|
|
'message': _("%s users queued; the background worker has started.", len(users)),
|
|
'type': 'success',
|
|
'next': {'type': 'ir.actions.act_window_close'},
|
|
},
|
|
}
|