# -*- coding: utf-8 -*- import logging from odoo import models, fields, api, _ _logger = logging.getLogger(__name__) class SocialMetaSyncQueue(models.Model): _name = 'social.meta.sync.queue' _description = 'Meta Catalog Sync Queue' _order = 'create_date desc, id desc' product_id = fields.Many2one('product.product', string='Product Variant', required=True, ondelete='cascade') product_tmpl_id = fields.Many2one('product.template', string='Product Template', related='product_id.product_tmpl_id', store=True) company_id = fields.Many2one('res.company', string='Company', required=True, default=lambda self: self.env.company) operation = fields.Selection([ ('create_update', 'Create / Update'), ('delete', 'Delete / Archive'), ], string='Operation', default='create_update', required=True) state = fields.Selection([ ('pending', 'Pending'), ('processing', 'Processing'), ('done', 'Synced'), ('failed', 'Failed'), ], string='Sync Status', default='pending', index=True, required=True) retry_count = fields.Integer(string='Retry Count', default=0) max_retries = fields.Integer(string='Max Retries', default=5) last_error = fields.Text(string='Last Error') last_attempt = fields.Datetime(string='Last Attempt') scheduled_at = fields.Datetime(string='Scheduled At', default=fields.Datetime.now) @api.model def enqueue_product(self, product, operation='create_update'): """Safely enqueues a product.product for synchronization without duplicating pending entries.""" if not product or not product.exists(): return False company = product.company_id or self.env.company if not company.meta_sync_enabled: return False existing = self.search([ ('product_id', '=', product.id), ('state', 'in', ('pending', 'processing')), ], limit=1) if existing: existing.write({ 'operation': operation, 'scheduled_at': fields.Datetime.now(), }) return existing return self.create({ 'product_id': product.id, 'company_id': company.id, 'operation': operation, 'state': 'pending', 'scheduled_at': fields.Datetime.now(), }) @api.model def process_queue(self, limit=50): """Processes pending items in the queue up to limit, grouping by company.""" records = self.search([ ('state', 'in', ('pending', 'failed')), ('retry_count', '<', 5), ('scheduled_at', '<=', fields.Datetime.now()), ], limit=limit) if not records: return 0 _logger.info("Processing Meta Sync Queue (%d items)...", len(records)) records.write({'state': 'processing'}) by_company = {} for rec in records: by_company.setdefault(rec.company_id, []).append(rec) meta_api = self.env['social.meta.catalog'] processed_count = 0 for company, queue_items in by_company.items(): if not company.meta_sync_enabled or not company.meta_catalog_id or not company.meta_access_token: for item in queue_items: item.write({ 'state': 'failed', 'last_error': _("Meta Catalog is not configured or enabled for this company."), 'last_attempt': fields.Datetime.now(), }) continue requests_list = [] item_map = {} for item in queue_items: product = item.product_id if not product.exists(): item.unlink() continue method = 'DELETE' if item.operation == 'delete' else 'UPDATE' try: payload = meta_api.build_item_payload(product, method=method) requests_list.append(payload) retailer_id = payload.get('retailer_id') item_map[retailer_id] = item except Exception as ex: _logger.exception("Failed building item payload for product %s", product.id) item.write({ 'state': 'failed', 'last_error': str(ex), 'retry_count': item.retry_count + 1, 'last_attempt': fields.Datetime.now(), }) if not requests_list: continue result = meta_api.sync_batch_items(company, requests_list) now = fields.Datetime.now() if result.get('success'): for item in item_map.values(): item.write({ 'state': 'done', 'last_error': False, 'last_attempt': now, }) item.product_id.sudo().write({ 'meta_sync_status': 'synced', 'meta_last_sync': now, 'meta_sync_error': False, 'meta_product_id': item.product_id.default_code or f"odoo_prod_{item.product_id.id}", }) if item.product_tmpl_id: item.product_tmpl_id.sudo().write({ 'meta_sync_status': 'synced', 'meta_last_sync': now, 'meta_sync_error': False, }) processed_count += len(item_map) else: err_msg = result.get('error', _("Unknown Meta API error")) for item in item_map.values(): new_retry = item.retry_count + 1 item.write({ 'state': 'failed' if new_retry >= item.max_retries else 'pending', 'retry_count': new_retry, 'last_error': err_msg, 'last_attempt': now, }) item.product_id.sudo().write({ 'meta_sync_status': 'failed', 'meta_last_sync': now, 'meta_sync_error': err_msg, }) return processed_count def action_retry(self): """Action for users to manually retry failed sync queue entries.""" for rec in self: rec.write({ 'state': 'pending', 'retry_count': 0, 'last_error': False, 'scheduled_at': fields.Datetime.now(), }) return True @api.model def action_sync_all_published_products(self, company=None): """Enqueues all published products for a full catalog refresh.""" company = company or self.env.company domain = [('is_published', '=', True), ('sale_ok', '=', True)] if company: domain.append(('company_id', 'in', (False, company.id))) templates = self.env['product.template'].search(domain) variants = templates.mapped('product_variant_ids') count = 0 for variant in variants: self.enqueue_product(variant, operation='create_update') count += 1 _logger.info("Enqueued %d published product variants for Meta catalog sync.", count) # Trigger immediate processing self.process_queue(limit=50) return { 'type': 'ir.actions.client', 'tag': 'display_notification', 'params': { 'title': _("Catalog Sync Enqueued"), 'message': _("%d published product variants have been queued for Meta catalog synchronization.") % count, 'type': 'info', 'sticky': False, } }