Upsert keyed entities on POST /api/<entity> (#1)
Test backoffice (smb-crm) / test (push) Has been cancelled

Easy!Appointments re-fires appointment_save on reschedule with the same
EA-<id> booking_id, so add_entity's plain INSERT 500s on the PK conflict.
Switch to ON CONFLICT (pk) DO UPDATE for every entity except activity_log
(no client-supplied pk), logging add vs update accordingly. Extended to
all keyed entities per the issue body, while preserving created_at/
received_at on redelivery for leads/clients/credentials so a redelivered
webhook can't clobber the original creation time.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
2026-08-04 13:42:00 +02:00
parent aabd1d56c4
commit 1b8f8c1391
2 changed files with 195 additions and 3 deletions
+35 -3
View File
@@ -129,6 +129,17 @@ def next_client_id(cur):
return "C-%04d" % (mx + 1) return "C-%04d" % (mx + 1)
# Columns that add_entity defaults to "now" only when the caller omits them
# (see the per-entity blocks below) -- a re-delivered webhook that omits the
# same column on a redelivery must not have it upserted back to a fresh
# "now", clobbering the original creation time of the row it's updating.
INSERT_ONLY_COLS = {
"leads": {"received_at"},
"clients": {"created_at"},
"credentials": {"created_at"},
}
@app.post("/api/<entity>") @app.post("/api/<entity>")
def add_entity(entity): def add_entity(entity):
if entity not in db.TABLES: if entity not in db.TABLES:
@@ -160,10 +171,31 @@ def add_entity(entity):
return jsonify({"error": f"{pk} required"}), 400 return jsonify({"error": f"{pk} required"}), 400
cols = spec["cols"] cols = spec["cols"]
ph = ", ".join(["%s"] * len(cols)) ph = ", ".join(["%s"] * len(cols))
cur.execute(f"INSERT INTO {entity} ({', '.join(cols)}) VALUES ({ph})", collist = ", ".join(cols)
[row.get(c) for c in cols]) if entity == "activity_log":
log_activity(cur, row.get("client_id"), f"add {entity}", f"{pk}={row.get(pk)}") # activity_log has no client-supplied primary key (its real pk,
# a serial "id", isn't even in cols) -- always a fresh row.
cur.execute(f"INSERT INTO {entity} ({collist}) VALUES ({ph})",
[row.get(c) for c in cols])
updated = False
else:
# Upsert by pk (#1): re-delivered webhooks -- e.g. Easy!Appointments
# re-firing appointment_save with the same EA-<id> booking_id on
# reschedule -- must update the existing row instead of 500ing on
# a PK conflict.
skip = INSERT_ONLY_COLS.get(entity, set()) | {pk}
updates = ", ".join(f"{c} = EXCLUDED.{c}" for c in cols if c not in skip)
cur.execute(
f"INSERT INTO {entity} ({collist}) VALUES ({ph}) "
f"ON CONFLICT ({pk}) DO UPDATE SET {updates} "
f"RETURNING (xmax = 0) AS inserted",
[row.get(c) for c in cols])
updated = not cur.fetchone()["inserted"]
log_activity(cur, row.get("client_id"),
f"{'update' if updated else 'add'} {entity}", f"{pk}={row.get(pk)}")
conn.commit() conn.commit()
if updated:
return jsonify({"updated": row.get(pk)}), 200
return jsonify({"added": row.get(pk)}), 201 return jsonify({"added": row.get(pk)}), 201
+160
View File
@@ -0,0 +1,160 @@
"""Flask test client / real-DB integration tests for the generic keyed-entity
POST endpoint (#1): add_entity must upsert on the primary key so re-delivered
webhooks (Easy!Appointments re-firing appointment_save on reschedule with the
same EA-<id> booking_id) update the existing row instead of 500ing on a PK
conflict.
"""
import db
import pytest
import app as app_module
from app import app as flask_app
CLIENT_A = "C-TEST-APIUP-A"
AUTH = {"X-CRM-Token": "test-crm-token"}
@pytest.fixture
def client(monkeypatch):
flask_app.config["TESTING"] = True
monkeypatch.setattr(app_module, "CRM_TOKEN", "test-crm-token")
return flask_app.test_client()
def _insert_client(client_id):
with db.connect() as conn, conn.cursor() as cur:
cur.execute(
"INSERT INTO clients (client_id, business_name) VALUES (%s, %s) "
"ON CONFLICT (client_id) DO NOTHING",
(client_id, "Café " + client_id))
conn.commit()
def _activity_rows(cur, client_id, action_prefix):
cur.execute(
"SELECT * FROM activity_log WHERE client_id = %s AND action LIKE %s ORDER BY ts",
(client_id, action_prefix + "%"))
return cur.fetchall()
# ---- bookings: reschedule re-delivery must upsert, not 500 ----
def test_reschedule_repost_updates_existing_booking(client):
_insert_client(CLIENT_A)
booking_id = "EA-42"
resp = client.post("/api/bookings", headers=AUTH, json={
"booking_id": booking_id, "client_id": CLIENT_A, "customer_name": "Ivy",
"service": "Haircut", "start_time": "2026-09-01T10:00:00Z",
"end_time": "2026-09-01T10:30:00Z", "source": "easyappointments",
"status": "confirmed"})
assert resp.status_code == 201
assert resp.get_json()["added"] == booking_id
# EA re-fires appointment_save on reschedule with the same booking_id.
resp2 = client.post("/api/bookings", headers=AUTH, json={
"booking_id": booking_id, "client_id": CLIENT_A, "customer_name": "Ivy",
"service": "Haircut", "start_time": "2026-09-01T11:00:00Z",
"end_time": "2026-09-01T11:30:00Z", "source": "easyappointments",
"status": "confirmed"})
assert resp2.status_code == 200
assert resp2.get_json()["updated"] == booking_id
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT * FROM bookings WHERE booking_id = %s", (booking_id,))
rows = cur.fetchall()
assert len(rows) == 1
assert rows[0]["start_time"].hour == 11
add_rows = _activity_rows(cur, CLIENT_A, "add bookings")
update_rows = _activity_rows(cur, CLIENT_A, "update bookings")
assert len(add_rows) == 1
assert len(update_rows) == 1
def test_cancel_repost_updates_status(client):
_insert_client(CLIENT_A)
booking_id = "EA-43"
client.post("/api/bookings", headers=AUTH, json={
"booking_id": booking_id, "client_id": CLIENT_A, "status": "confirmed"})
resp = client.post("/api/bookings", headers=AUTH, json={
"booking_id": booking_id, "client_id": CLIENT_A, "status": "cancelled"})
assert resp.status_code == 200
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT status FROM bookings WHERE booking_id = %s", (booking_id,))
assert cur.fetchone()["status"] == "cancelled"
# ---- other keyed entities: same upsert treatment ----
def test_repost_same_invoice_id_updates_instead_of_500(client):
_insert_client(CLIENT_A)
invoice_id = "INV-TEST-1"
resp = client.post("/api/invoices", headers=AUTH, json={
"invoice_id": invoice_id, "client_id": CLIENT_A, "amount_eur": "100",
"status": "open"})
assert resp.status_code == 201
resp2 = client.post("/api/invoices", headers=AUTH, json={
"invoice_id": invoice_id, "client_id": CLIENT_A, "amount_eur": "100",
"status": "paid"})
assert resp2.status_code == 200
assert resp2.get_json()["updated"] == invoice_id
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT status FROM invoices WHERE invoice_id = %s", (invoice_id,))
assert cur.fetchone()["status"] == "paid"
def test_repost_same_lead_id_preserves_received_at(client):
lead_id = "L-TEST-1"
resp = client.post("/api/leads", headers=AUTH, json={
"lead_id": lead_id, "client_id": CLIENT_A, "name": "Ivy", "status": "new"})
assert resp.status_code == 201
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT received_at FROM leads WHERE lead_id = %s", (lead_id,))
first_received_at = cur.fetchone()["received_at"]
# Re-delivered webhook omits received_at -- must not clobber the original.
resp2 = client.post("/api/leads", headers=AUTH, json={
"lead_id": lead_id, "client_id": CLIENT_A, "name": "Ivy", "status": "contacted"})
assert resp2.status_code == 200
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT received_at, status FROM leads WHERE lead_id = %s", (lead_id,))
row = cur.fetchone()
assert row["received_at"] == first_received_at
assert row["status"] == "contacted"
def test_repost_same_client_id_preserves_created_at(client):
client_id = "C-TEST-APIUP-CREATED"
resp = client.post("/api/clients", headers=AUTH, json={
"client_id": client_id, "business_name": "Café Redelivered"})
assert resp.status_code == 201
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT created_at FROM clients WHERE client_id = %s", (client_id,))
first_created_at = cur.fetchone()["created_at"]
resp2 = client.post("/api/clients", headers=AUTH, json={
"client_id": client_id, "business_name": "Café Redelivered", "status": "active"})
assert resp2.status_code == 200
with db.connect() as conn, conn.cursor() as cur:
cur.execute("SELECT created_at, status FROM clients WHERE client_id = %s", (client_id,))
row = cur.fetchone()
assert row["created_at"] == first_created_at
assert row["status"] == "active"
# ---- activity_log has no client-supplied primary key: plain insert only ----
def test_activity_log_post_is_plain_insert_not_upsert(client):
resp1 = client.post("/api/activity_log", headers=AUTH, json={
"ts": "2026-09-01T00:00:00Z", "workflow": "test", "client_id": CLIENT_A,
"action": "note", "detail": "one"})
resp2 = client.post("/api/activity_log", headers=AUTH, json={
"ts": "2026-09-01T00:00:00Z", "workflow": "test", "client_id": CLIENT_A,
"action": "note", "detail": "two"})
assert resp1.status_code == 201
assert resp2.status_code == 201
with db.connect() as conn, conn.cursor() as cur:
cur.execute(
"SELECT count(*) AS n FROM activity_log WHERE client_id = %s AND action = 'note'",
(CLIENT_A,))
assert cur.fetchone()["n"] == 2