feat: add POST /sync endpoint + daily n8n drift check workflow
Adds extension sync that queries each customer's EspoCRM to compare actually-installed extensions with what the platform DB tracks. Updates platform DB on version mismatches and reports drifts to Mattermost. n8n workflow ID: NXxBHm3qhWL3yHza (daily at 07:00) Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,169 @@
|
||||
{
|
||||
"name": "Extension Sync: Daily Drift Check",
|
||||
"nodes": [
|
||||
{
|
||||
"parameters": {
|
||||
"rule": {
|
||||
"interval": [
|
||||
{
|
||||
"triggerAtHour": 7
|
||||
}
|
||||
]
|
||||
}
|
||||
},
|
||||
"id": "schedule-trigger",
|
||||
"name": "Daily 07:00",
|
||||
"type": "n8n-nodes-base.scheduleTrigger",
|
||||
"typeVersion": 1.2,
|
||||
"position": [0, 0]
|
||||
},
|
||||
{
|
||||
"parameters": {
|
||||
"method": "POST",
|
||||
"url": "https://platform.dev.marcus-law.co.il/sync",
|
||||
"sendBody": true,
|
||||
"bodyParameters": {
|
||||
"parameters": []
|
||||
},
|
||||
"options": {
|
||||
"timeout": 120000
|
||||
}
|
||||
},
|
||||
"id": "call-sync-api",
|
||||
"name": "Call Platform Sync",
|
||||
"type": "n8n-nodes-base.httpRequest",
|
||||
"typeVersion": 4.2,
|
||||
"position": [220, 0]
|
||||
},
|
||||
{
|
||||
"parameters": {
|
||||
"conditions": {
|
||||
"options": {
|
||||
"caseSensitive": true,
|
||||
"leftValue": "",
|
||||
"typeValidation": "strict"
|
||||
},
|
||||
"conditions": [
|
||||
{
|
||||
"id": "drift-check",
|
||||
"leftValue": "={{ $json.total_drifts }}",
|
||||
"rightValue": 0,
|
||||
"operator": {
|
||||
"type": "number",
|
||||
"operation": "gt"
|
||||
}
|
||||
}
|
||||
],
|
||||
"combinator": "and"
|
||||
},
|
||||
"options": {}
|
||||
},
|
||||
"id": "check-drifts",
|
||||
"name": "Has Drifts?",
|
||||
"type": "n8n-nodes-base.if",
|
||||
"typeVersion": 2,
|
||||
"position": [440, 0]
|
||||
},
|
||||
{
|
||||
"parameters": {
|
||||
"jsCode": "// Build Mattermost message from sync results\nconst syncResult = $input.first().json;\nconst results = syncResult.results || [];\n\nlet lines = [\n `### :warning: Extension Drift Report`,\n `**${syncResult.customers_checked}** customers checked | **${syncResult.total_drifts}** drifts found`,\n ''\n];\n\nfor (const r of results) {\n if (r.status === 'error') {\n lines.push(`#### :x: ${r.customer}`);\n lines.push(`> ${r.message}`);\n lines.push('');\n continue;\n }\n if (!r.drifts || r.drifts.length === 0) continue;\n\n lines.push(`#### :mag: ${r.customer} (${r.drift_count} drifts)`);\n lines.push(`| Extension | Type | Platform | Server |`);\n lines.push(`|-----------|------|----------|--------|`);\n\n for (const d of r.drifts) {\n const type = {\n 'version_mismatch': ':arrows_counterclockwise: Version mismatch',\n 'missing_on_server': ':red_circle: Missing on server',\n 'extra_on_server': ':large_blue_circle: Extra on server'\n }[d.type] || d.type;\n lines.push(`| ${d.extension} | ${type} | ${d.platform || '-'} | ${d.server || '-'} |`);\n }\n\n if (r.updated > 0) {\n lines.push(`> :white_check_mark: Updated ${r.updated} records in platform DB`);\n }\n lines.push('');\n}\n\nreturn [{ json: { message: lines.join('\\n') } }];"
|
||||
},
|
||||
"id": "format-drift-message",
|
||||
"name": "Format Drift Report",
|
||||
"type": "n8n-nodes-base.code",
|
||||
"typeVersion": 2,
|
||||
"position": [660, -100]
|
||||
},
|
||||
{
|
||||
"parameters": {
|
||||
"jsCode": "// Build success summary\nconst syncResult = $input.first().json;\n\nconst message = [\n `### :white_check_mark: Extension Sync — All Clear`,\n `**${syncResult.customers_checked}** customers checked, no drifts found.`,\n '',\n `| Customer | Server | Platform |`,\n `|----------|--------|----------|`\n];\n\nfor (const r of (syncResult.results || [])) {\n if (r.status === 'error') {\n message.push(`| ${r.customer} | :x: Error | - |`);\n } else {\n message.push(`| ${r.customer} | ${r.server_count} exts | ${r.platform_count} tracked |`);\n }\n}\n\nreturn [{ json: { message: message.join('\\n') } }];"
|
||||
},
|
||||
"id": "format-ok-message",
|
||||
"name": "Format OK Summary",
|
||||
"type": "n8n-nodes-base.code",
|
||||
"typeVersion": 2,
|
||||
"position": [660, 100]
|
||||
},
|
||||
{
|
||||
"parameters": {
|
||||
"method": "POST",
|
||||
"url": "https://mattermost.dev.marcus-law.co.il/api/v4/posts",
|
||||
"authentication": "genericCredentialType",
|
||||
"genericAuthType": "httpHeaderAuth",
|
||||
"sendBody": true,
|
||||
"specifyBody": "json",
|
||||
"jsonBody": "={{ JSON.stringify({ channel_id: '5x1mzsm7djdz7ydqkde93dkq4e', message: $json.message }) }}",
|
||||
"options": {}
|
||||
},
|
||||
"id": "mattermost-drift",
|
||||
"name": "Mattermost (Drift)",
|
||||
"type": "n8n-nodes-base.httpRequest",
|
||||
"typeVersion": 4.2,
|
||||
"position": [880, -100],
|
||||
"credentials": {
|
||||
"httpHeaderAuth": {
|
||||
"id": "mattermost-token",
|
||||
"name": "Mattermost Token"
|
||||
}
|
||||
}
|
||||
},
|
||||
{
|
||||
"parameters": {
|
||||
"method": "POST",
|
||||
"url": "https://mattermost.dev.marcus-law.co.il/api/v4/posts",
|
||||
"authentication": "genericCredentialType",
|
||||
"genericAuthType": "httpHeaderAuth",
|
||||
"sendBody": true,
|
||||
"specifyBody": "json",
|
||||
"jsonBody": "={{ JSON.stringify({ channel_id: '5x1mzsm7djdz7ydqkde93dkq4e', message: $json.message }) }}",
|
||||
"options": {}
|
||||
},
|
||||
"id": "mattermost-ok",
|
||||
"name": "Mattermost (OK)",
|
||||
"type": "n8n-nodes-base.httpRequest",
|
||||
"typeVersion": 4.2,
|
||||
"position": [880, 100],
|
||||
"credentials": {
|
||||
"httpHeaderAuth": {
|
||||
"id": "mattermost-token",
|
||||
"name": "Mattermost Token"
|
||||
}
|
||||
}
|
||||
}
|
||||
],
|
||||
"connections": {
|
||||
"Daily 07:00": {
|
||||
"main": [
|
||||
[{ "node": "Call Platform Sync", "type": "main", "index": 0 }]
|
||||
]
|
||||
},
|
||||
"Call Platform Sync": {
|
||||
"main": [
|
||||
[{ "node": "Has Drifts?", "type": "main", "index": 0 }]
|
||||
]
|
||||
},
|
||||
"Has Drifts?": {
|
||||
"main": [
|
||||
[{ "node": "Format Drift Report", "type": "main", "index": 0 }],
|
||||
[{ "node": "Format OK Summary", "type": "main", "index": 0 }]
|
||||
]
|
||||
},
|
||||
"Format Drift Report": {
|
||||
"main": [
|
||||
[{ "node": "Mattermost (Drift)", "type": "main", "index": 0 }]
|
||||
]
|
||||
},
|
||||
"Format OK Summary": {
|
||||
"main": [
|
||||
[{ "node": "Mattermost (OK)", "type": "main", "index": 0 }]
|
||||
]
|
||||
}
|
||||
},
|
||||
"settings": {
|
||||
"executionOrder": "v1"
|
||||
},
|
||||
"tags": [
|
||||
{ "name": "extension-platform" },
|
||||
{ "name": "daily" }
|
||||
]
|
||||
}
|
||||
@@ -8,6 +8,7 @@ Pipeline endpoints (called by n8n):
|
||||
POST /publish — Download from Gitea, build ZIP, publish to Composer Registry + Release
|
||||
POST /deploy — Deploy extension to a specific customer via EspoCRM API
|
||||
POST /deploy-all — Deploy to all licensed customers
|
||||
POST /sync — Sync installed extensions from EspoCRM servers with platform DB
|
||||
|
||||
Dashboard API:
|
||||
GET /api/stats
|
||||
@@ -357,6 +358,81 @@ def deploy_to_customer(customer, module, version):
|
||||
return {'status': 'ok', 'ext_id': ext_id}
|
||||
|
||||
|
||||
# --- Sync logic ---
|
||||
def sync_customer_extensions(customer):
|
||||
"""Query a customer's EspoCRM to see what's really installed, compare with platform DB."""
|
||||
espo_url = customer['espocrm_url'].rstrip('/')
|
||||
api_key = customer['espocrm_api_key']
|
||||
cid = customer['id']
|
||||
|
||||
# 1. Query EspoCRM for installed extensions
|
||||
headers = {'Espo-Authorization': api_key}
|
||||
req = urllib.request.Request(f'{espo_url}/api/v1/Extension?select=name,version&maxSize=200',
|
||||
headers=headers)
|
||||
try:
|
||||
resp = urllib.request.urlopen(req, context=SSL_CTX, timeout=30)
|
||||
data = json.loads(resp.read())
|
||||
server_exts = {e['name']: e.get('version', '') for e in data.get('list', [])}
|
||||
except Exception as e:
|
||||
return {'customer': customer['name'], 'status': 'error',
|
||||
'message': f'Cannot reach EspoCRM: {str(e)[:200]}'}
|
||||
|
||||
# 2. Get what platform thinks is installed
|
||||
platform_exts = query("""
|
||||
SELECT e.name, ce.installed_version, ce.license_status
|
||||
FROM customer_extensions ce JOIN extensions e ON ce.extension_id = e.id
|
||||
WHERE ce.customer_id = %s AND ce.license_status = 'active'""", (cid,))
|
||||
platform_map = {r['name']: r['installed_version'] for r in platform_exts}
|
||||
|
||||
# 3. Compare
|
||||
drifts = []
|
||||
updated = 0
|
||||
|
||||
# Check each platform-tracked extension
|
||||
for ext_name, platform_ver in platform_map.items():
|
||||
server_ver = server_exts.get(ext_name)
|
||||
if server_ver is None:
|
||||
drifts.append({'extension': ext_name, 'type': 'missing_on_server',
|
||||
'platform': platform_ver, 'server': None})
|
||||
elif server_ver != platform_ver:
|
||||
drifts.append({'extension': ext_name, 'type': 'version_mismatch',
|
||||
'platform': platform_ver, 'server': server_ver})
|
||||
# Update platform DB to reflect reality
|
||||
query("""UPDATE customer_extensions SET installed_version = %s, updated_at = NOW()
|
||||
FROM extensions e WHERE customer_extensions.extension_id = e.id
|
||||
AND customer_extensions.customer_id = %s AND e.name = %s""",
|
||||
(server_ver, cid, ext_name), fetch=None)
|
||||
updated += 1
|
||||
|
||||
# Check for extensions on server not tracked by platform
|
||||
for ext_name, server_ver in server_exts.items():
|
||||
if ext_name not in platform_map and ext_name in MODULES:
|
||||
drifts.append({'extension': ext_name, 'type': 'extra_on_server',
|
||||
'platform': None, 'server': server_ver})
|
||||
|
||||
return {
|
||||
'customer': customer['name'], 'customer_id': cid, 'status': 'ok',
|
||||
'server_count': len(server_exts), 'platform_count': len(platform_map),
|
||||
'drifts': drifts, 'drift_count': len(drifts), 'updated': updated
|
||||
}
|
||||
|
||||
|
||||
def sync_all_customers():
|
||||
"""Run sync for all active customers."""
|
||||
customers = query("""SELECT id, name, espocrm_url, espocrm_api_key
|
||||
FROM customers WHERE is_active AND espocrm_api_key IS NOT NULL""")
|
||||
results = []
|
||||
total_drifts = 0
|
||||
for c in customers:
|
||||
r = sync_customer_extensions(c)
|
||||
results.append(r)
|
||||
total_drifts += r.get('drift_count', 0)
|
||||
return {
|
||||
'status': 'ok', 'customers_checked': len(results),
|
||||
'total_drifts': total_drifts, 'results': results
|
||||
}
|
||||
|
||||
|
||||
# --- HTTP Handler ---
|
||||
class APIHandler(BaseHTTPRequestHandler):
|
||||
|
||||
@@ -392,6 +468,8 @@ class APIHandler(BaseHTTPRequestHandler):
|
||||
self._assign_extension(int(re.search(r'/api/customers/(\d+)/extensions', path).group(1)), body)
|
||||
elif path == '/api/deploy':
|
||||
self._handle_deploy(body)
|
||||
elif path == '/sync':
|
||||
self._handle_sync(body)
|
||||
else:
|
||||
self._respond(404, {'error': 'Not found'})
|
||||
|
||||
@@ -544,6 +622,19 @@ class APIHandler(BaseHTTPRequestHandler):
|
||||
self._respond(200, {'status': 'ok', 'module': module, 'version': version,
|
||||
'deployed_to': len(results), 'results': results})
|
||||
|
||||
# --- Sync ---
|
||||
def _handle_sync(self, body):
|
||||
cid = body.get('customer_id')
|
||||
if cid:
|
||||
customer = query("SELECT id, name, espocrm_url, espocrm_api_key FROM customers WHERE id = %s AND is_active", (cid,), fetch='one')
|
||||
if not customer:
|
||||
self._respond(404, {'error': 'Customer not found'}); return
|
||||
result = sync_customer_extensions(customer)
|
||||
self._respond(200, result)
|
||||
else:
|
||||
result = sync_all_customers()
|
||||
self._respond(200, result)
|
||||
|
||||
# --- Helpers ---
|
||||
def _read_body(self):
|
||||
length = int(self.headers.get('Content-Length', 0))
|
||||
|
||||
Reference in New Issue
Block a user