diff --git a/db/migrations/0011_suppress_shopify_legacy_outbox_events.sql b/db/migrations/0011_suppress_shopify_legacy_outbox_events.sql new file mode 100644 index 0000000..072ce80 --- /dev/null +++ b/db/migrations/0011_suppress_shopify_legacy_outbox_events.sql @@ -0,0 +1,110 @@ +BEGIN; + +-- Shopify delta imports are authoritative reads from Shopify. They must not +-- enter the legacy n8n outbox path; future operational Shopify events need an +-- explicit integration contract instead. +CREATE OR REPLACE FUNCTION fn_enqueue_sales_order_events() +RETURNS TRIGGER +LANGUAGE plpgsql +AS $$ +DECLARE + v_event_key TEXT; +BEGIN + IF NEW.order_source = 'shopify' THEN + RETURN NEW; + END IF; + + IF TG_OP = 'INSERT' THEN + v_event_key := format( + 'order.imported:%s:%s', + NEW.external_ref, + to_char(NEW.imported_at, 'YYYYMMDDHH24MISSUS') + ); + PERFORM fn_enqueue_event( + 'order.imported', v_event_key, 'sales_order', NEW.id::TEXT, + jsonb_build_object( + 'orderId', NEW.id, 'externalRef', NEW.external_ref, + 'orderSource', NEW.order_source, 'orderStatus', NEW.order_status, + 'paymentStatus', NEW.payment_status, 'occurredAt', NOW() + ) + ); + ELSIF TG_OP = 'UPDATE' THEN + IF OLD.imported_at IS DISTINCT FROM NEW.imported_at THEN + v_event_key := format( + 'order.imported:%s:%s', NEW.external_ref, + to_char(NEW.imported_at, 'YYYYMMDDHH24MISSUS') + ); + PERFORM fn_enqueue_event( + 'order.imported', v_event_key, 'sales_order', NEW.id::TEXT, + jsonb_build_object( + 'orderId', NEW.id, 'externalRef', NEW.external_ref, + 'orderSource', NEW.order_source, 'orderStatus', NEW.order_status, + 'paymentStatus', NEW.payment_status, 'occurredAt', NOW() + ) + ); + END IF; + + IF OLD.order_status IS DISTINCT FROM NEW.order_status + AND NEW.order_status = 'cancelled' THEN + v_event_key := format( + 'order.cancelled.full:%s:%s', NEW.external_ref, + COALESCE(NEW.cancelled_at, NOW()) + ); + PERFORM fn_enqueue_event( + 'order.cancelled.full', v_event_key, 'sales_order', NEW.id::TEXT, + jsonb_build_object( + 'orderId', NEW.id, 'externalRef', NEW.external_ref, + 'orderSource', NEW.order_source, 'orderStatus', NEW.order_status, + 'cancelledAt', COALESCE(NEW.cancelled_at, NOW()), + 'cancelledReason', NEW.cancelled_reason + ) + ); + END IF; + END IF; + + RETURN NEW; +END; +$$; + +CREATE OR REPLACE FUNCTION fn_enqueue_partial_cancel_events() +RETURNS TRIGGER +LANGUAGE plpgsql +AS $$ +DECLARE + v_order_id BIGINT; + v_external_ref TEXT; + v_order_source TEXT; + v_event_key TEXT; +BEGIN + IF NEW.line_status = 'partially_cancelled' + AND OLD.line_status IS DISTINCT FROM NEW.line_status THEN + SELECT so.id, so.external_ref, so.order_source + INTO v_order_id, v_external_ref, v_order_source + FROM sales_order so + WHERE so.id = NEW.sales_order_id; + + IF v_order_source = 'shopify' THEN + RETURN NEW; + END IF; + + v_event_key := format( + 'order.cancelled.partial:%s:%s:%s', + v_external_ref, NEW.id, NEW.qty_cancelled + ); + PERFORM fn_enqueue_event( + 'order.cancelled.partial', v_event_key, 'sales_order_line', NEW.id::TEXT, + jsonb_build_object( + 'orderId', v_order_id, 'externalRef', v_external_ref, + 'orderSource', v_order_source, 'lineId', NEW.id, + 'lineNo', NEW.line_no, 'qty', NEW.qty, + 'qtyCancelled', NEW.qty_cancelled, 'lineStatus', NEW.line_status, + 'occurredAt', NOW() + ) + ); + END IF; + + RETURN NEW; +END; +$$; + +COMMIT; diff --git a/modules/erp/bestellungen/shopify-projection.php b/modules/erp/bestellungen/shopify-projection.php new file mode 100644 index 0000000..d737dc2 --- /dev/null +++ b/modules/erp/bestellungen/shopify-projection.php @@ -0,0 +1,306 @@ + 'paid', + 'AUTHORIZED' => 'authorized', + 'PARTIALLY_REFUNDED' => 'partially_refunded', + 'REFUNDED' => 'refunded', + 'VOIDED' => 'voided', + default => 'unknown', + }; +} + +function shopify_projection_order_status(array $order): string +{ + if (($order['cancelledAt'] ?? null) !== null) { + return 'cancelled'; + } + + return strtoupper((string) ($order['displayFulfillmentStatus'] ?? '')) === 'FULFILLED' + ? 'fulfilled' + : 'imported'; +} + +function shopify_projection_money(array $order, string $field): ?float +{ + $amount = $order[$field]['shopMoney']['amount'] ?? null; + return $amount === null ? null : parse_number($amount); +} + +function shopify_projection_customer(PDO $pdo, array $customer, bool $dryRun): ?int +{ + $customerGid = trim((string) ($customer['id'] ?? '')); + if ($customerGid === '') { + return null; + } + + $identity = $pdo->prepare( + 'SELECT party_id FROM party_external_identity + WHERE source_system = \'shopify\' AND external_id = :external_id LIMIT 1' + ); + $identity->execute([':external_id' => $customerGid]); + $partyId = $identity->fetchColumn(); + if ($partyId !== false) { + return (int) $partyId; + } + + $email = trim((string) ($customer['email'] ?? '')); + if ($dryRun) { + return null; + } + + $name = trim((string) ($customer['firstName'] ?? '') . ' ' . (string) ($customer['lastName'] ?? '')); + if ($name === '') { + $name = 'Shopify-Kunde'; + } + + $partyStmt = $pdo->prepare( + "INSERT INTO party (type, name, email, status, created_at, updated_at) + VALUES ('customer', :name, :email, 'active', NOW(), NOW()) + RETURNING id" + ); + $partyStmt->execute([':name' => $name, ':email' => $email !== '' ? $email : null]); + $partyId = $partyStmt->fetchColumn(); + if ($partyId === false) { + throw new RuntimeException('Shopify-Kunde konnte nicht angelegt werden'); + } + + $identityStmt = $pdo->prepare( + "INSERT INTO party_external_identity (party_id, source_system, external_id, observed_at, created_at, updated_at) + VALUES (:party_id, 'shopify', :external_id, NOW(), NOW(), NOW())" + ); + $identityStmt->execute([':party_id' => (int) $partyId, ':external_id' => $customerGid]); + + return (int) $partyId; +} + +function shopify_projection_address(PDO $pdo, ?int $partyId, array $address, string $type, bool $dryRun): void +{ + if ($partyId === null || $dryRun || $address === []) { + return; + } + + $stmt = $pdo->prepare( + 'INSERT INTO address ( + party_id, type, first_name, last_name, company_name, street, house_number, + zip, city, state_code, country_name, country_iso2, created_at, updated_at + ) VALUES ( + :party_id, :type, :first_name, :last_name, :company_name, :street, :house_number, + :zip, :city, :state_code, :country_name, :country_iso2, NOW(), NOW() + )' + ); + $street = trim((string) ($address['address1'] ?? '')); + $address2 = trim((string) ($address['address2'] ?? '')); + if ($address2 !== '') { + $street .= ($street !== '' ? ', ' : '') . $address2; + } + $stmt->execute([ + ':party_id' => $partyId, + ':type' => $type, + ':first_name' => trim((string) ($address['firstName'] ?? '')), + ':last_name' => trim((string) ($address['lastName'] ?? '')), + ':company_name' => trim((string) ($address['company'] ?? '')) ?: null, + ':street' => $street, + ':house_number' => null, + ':zip' => trim((string) ($address['zip'] ?? '')), + ':city' => trim((string) ($address['city'] ?? '')), + ':state_code' => trim((string) ($address['provinceCode'] ?? '')) ?: null, + ':country_name' => trim((string) ($address['country'] ?? '')) ?: null, + ':country_iso2' => trim((string) ($address['countryCodeV2'] ?? '')) ?: null, + ]); +} + +function shopify_projection_line(PDO $pdo, array $line, bool $dryRun): array +{ + $sku = trim((string) ($line['sku'] ?? ($line['variant']['sku'] ?? ''))); + $variantGid = trim((string) ($line['variant']['id'] ?? '')); + $productGid = trim((string) ($line['variant']['product']['id'] ?? '')); + $sellableItemId = find_shopify_sellable_item_id($pdo, $sku, $variantGid); + if ($sellableItemId === null) { + throw new RuntimeException("Kein Shopify-Artikel-Mapping fuer SKU '{$sku}'"); + } + + if (!$dryRun) { + bind_shopify_item_identity($pdo, $sellableItemId, $sku, $productGid, $variantGid); + } + + $qty = parse_number($line['quantity'] ?? null); + if ($qty === null || $qty <= 0) { + throw new RuntimeException("Ungueltige Shopify-Menge fuer SKU '{$sku}'"); + } + + return [ + 'line' => $line, + 'sku' => $sku, + 'variant_gid' => $variantGid, + 'product_gid' => $productGid, + 'sellable_item_id' => $sellableItemId, + 'qty' => $qty, + 'unit_price' => parse_number($line['originalUnitPriceSet']['shopMoney']['amount'] ?? null), + ]; +} + +function shopify_projection_stock_check(PDO $pdo, array $resolvedLines): void +{ + foreach ($resolvedLines as $resolvedLine) { + $components = get_item_components($pdo, (int) $resolvedLine['sellable_item_id']); + if ($components === []) { + throw new RuntimeException("Keine Komponenten fuer SKU '{$resolvedLine['sku']}'"); + } + foreach ($components as $component) { + $required = (float) $resolvedLine['qty'] * (float) $component['qty_per_item']; + $stmt = $pdo->prepare( + "SELECT COALESCE(SUM(v.qty_net), 0) + FROM v_stock_lot_balance v + JOIN stock_lot sl ON sl.id = v.stock_lot_id + WHERE sl.product_id = :product_id AND sl.status = 'current'" + ); + $stmt->execute([':product_id' => (int) $component['product_id']]); + $available = (float) ($stmt->fetchColumn() ?: 0); + if ($available + 0.0000001 < $required) { + throw new RuntimeException("Nicht genuegend aktueller Bestand fuer SKU '{$resolvedLine['sku']}'"); + } + } + } +} + +function project_shopify_order(PDO $pdo, array $order, array $env, bool $dryRun = true): array +{ + $orderGid = trim((string) ($order['id'] ?? '')); + if ($orderGid === '') { + throw new RuntimeException('Shopify Order-GID fehlt'); + } + + $existingStmt = $pdo->prepare('SELECT id FROM sales_order WHERE shopify_order_gid = :shopify_order_gid LIMIT 1'); + $existingStmt->execute([':shopify_order_gid' => $orderGid]); + $existingId = $existingStmt->fetchColumn(); + if ($existingId !== false) { + return ['status' => 'already_imported', 'shopify_order_gid' => $orderGid, 'sales_order_id' => (int) $existingId]; + } + + $resolvedLines = []; + foreach (($order['lineItems']['nodes'] ?? []) as $line) { + $resolvedLines[] = shopify_projection_line($pdo, $line, $dryRun); + } + if ($resolvedLines === []) { + throw new RuntimeException('Shopify-Bestellung enthaelt keine Positionen'); + } + shopify_projection_stock_check($pdo, $resolvedLines); + + if ($dryRun) { + $customer = is_array($order['customer'] ?? null) ? $order['customer'] : []; + shopify_projection_customer($pdo, $customer, true); + return [ + 'status' => 'dry_run_ready', + 'shopify_order_gid' => $orderGid, + 'sales_order_id' => null, + 'line_count' => count($resolvedLines), + 'allocation_count' => 0, + ]; + } + + $pdo->beginTransaction(); + try { + $customer = is_array($order['customer'] ?? null) ? $order['customer'] : []; + $partyId = shopify_projection_customer($pdo, $customer, false); + $orderStatus = shopify_projection_order_status($order); + $paymentStatus = shopify_projection_payment_status((string) ($order['displayFinancialStatus'] ?? '')); + $orderDate = new DateTimeImmutable((string) $order['createdAt']); + $externalRef = (string) ($order['name'] ?? $orderGid); + $orderStmt = $pdo->prepare( + 'INSERT INTO sales_order ( + external_ref, party_id, order_date, order_status, payment_status, amount_net, + amount_shipping, amount_tax, total_amount, currency, webhook_payload, imported_at, + created_at, updated_at, order_source, shopify_order_gid, shopify_order_name, + shopify_created_at, shopify_updated_at, shopify_financial_status, shopify_fulfillment_status + ) VALUES ( + :external_ref, :party_id, :order_date, :order_status, :payment_status, :amount_net, + :amount_shipping, :amount_tax, :total_amount, :currency, :webhook_payload::jsonb, NOW(), + NOW(), NOW(), \'shopify\', :shopify_order_gid, :shopify_order_name, + :shopify_created_at, :shopify_updated_at, :shopify_financial_status, :shopify_fulfillment_status + ) RETURNING id' + ); + $orderStmt->execute([ + ':external_ref' => $externalRef, + ':party_id' => $partyId, + ':order_date' => $orderDate->format('Y-m-d H:i:s'), + ':order_status' => $orderStatus, + ':payment_status' => $paymentStatus, + ':amount_net' => shopify_projection_money($order, 'currentSubtotalPriceSet'), + ':amount_shipping' => shopify_projection_money($order, 'totalShippingPriceSet'), + ':amount_tax' => shopify_projection_money($order, 'totalTaxSet'), + ':total_amount' => shopify_projection_money($order, 'totalPriceSet'), + ':currency' => (string) (($order['totalPriceSet']['shopMoney']['currencyCode'] ?? 'CHF')), + ':webhook_payload' => json_encode($order, JSON_THROW_ON_ERROR | JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES), + ':shopify_order_gid' => $orderGid, + ':shopify_order_name' => (string) ($order['name'] ?? ''), + ':shopify_created_at' => $orderDate->format('Y-m-d H:i:s'), + ':shopify_updated_at' => (new DateTimeImmutable((string) $order['updatedAt']))->format('Y-m-d H:i:s'), + ':shopify_financial_status' => (string) ($order['displayFinancialStatus'] ?? ''), + ':shopify_fulfillment_status' => (string) ($order['displayFulfillmentStatus'] ?? ''), + ]); + $orderId = (int) $orderStmt->fetchColumn(); + + $locations = get_default_location_ids($pdo); + $lineNo = 0; + $allocationCount = 0; + foreach ($resolvedLines as $resolvedLine) { + $lineNo++; + $line = $resolvedLine['line']; + $lineStmt = $pdo->prepare( + 'INSERT INTO sales_order_line ( + sales_order_id, line_no, sellable_item_id, raw_external_article_number, + raw_external_title, qty, unit_price, line_total, shopify_line_item_gid, + shopify_variant_gid, shopify_product_gid, sku, created_at, updated_at + ) VALUES ( + :sales_order_id, :line_no, :sellable_item_id, :article_number, + :title, :qty, :unit_price, :line_total, :line_gid, + :variant_gid, :product_gid, :sku, NOW(), NOW() + ) RETURNING id' + ); + $qty = (float) $resolvedLine['qty']; + $unitPrice = $resolvedLine['unit_price']; + $lineStmt->execute([ + ':sales_order_id' => $orderId, + ':line_no' => $lineNo, + ':sellable_item_id' => $resolvedLine['sellable_item_id'], + ':article_number' => $resolvedLine['sku'], + ':title' => (string) ($line['name'] ?? $resolvedLine['sku']), + ':qty' => $qty, + ':unit_price' => $unitPrice, + ':line_total' => $unitPrice === null ? null : round($qty * $unitPrice, 2), + ':line_gid' => $resolvedLine['line']['id'], + ':variant_gid' => $resolvedLine['variant_gid'], + ':product_gid' => $resolvedLine['product_gid'], + ':sku' => $resolvedLine['sku'], + ]); + $lineId = (int) $lineStmt->fetchColumn(); + $allocation = allocate_line_inventory($pdo, $orderId, $lineId, $lineNo, $qty, (int) $resolvedLine['sellable_item_id'], $locations, 'shopify-order'); + $allocationCount += (int) ($allocation['allocationCount'] ?? 0); + } + + shopify_projection_address($pdo, $partyId, (array) ($order['billingAddress'] ?? []), 'billing', false); + shopify_projection_address($pdo, $partyId, (array) ($order['shippingAddress'] ?? []), 'shipping', false); + $pdo->commit(); + } catch (Throwable $exception) { + if ($pdo->inTransaction()) { + $pdo->rollBack(); + } + throw $exception; + } + + return [ + 'status' => 'imported', + 'shopify_order_gid' => $orderGid, + 'sales_order_id' => $orderId, + 'line_count' => count($resolvedLines), + 'allocation_count' => $allocationCount, + ]; +} diff --git a/modules/erp/import-integration/shopify-delta-run.php b/modules/erp/import-integration/shopify-delta-run.php new file mode 100644 index 0000000..8d2763d --- /dev/null +++ b/modules/erp/import-integration/shopify-delta-run.php @@ -0,0 +1,169 @@ +prepare( + "INSERT INTO process_runs (process_name, status, scope_json, requested_by) + VALUES ('erp.import-integration.shopify_delta_order_import', 'running', :scope_json::jsonb, 'cli') + RETURNING id" + ); + $stmt->execute([':scope_json' => json_encode([ + 'cutoff_timestamp' => $cutoff, + 'dry_run' => $dryRun, + ], JSON_THROW_ON_ERROR)]); + return (int) $stmt->fetchColumn(); +} + +function shopify_delta_finish_run(PDO $pdo, int $runId, string $status, array $result, ?Throwable $error = null): void +{ + $stmt = $pdo->prepare( + 'UPDATE process_runs + SET status = :status, result_json = :result_json::jsonb, finished_at = NOW(), + error_code = :error_code, error_message = :error_message + WHERE id = :id' + ); + $stmt->execute([ + ':status' => $status, + ':result_json' => json_encode($result, JSON_THROW_ON_ERROR), + ':error_code' => $error === null ? null : 'shopify_delta_failed', + ':error_message' => $error?->getMessage(), + ':id' => $runId, + ]); +} + +function shopify_delta_load_orders(array $env, string $searchQuery): array +{ + $query = <<<'GQL' +query ($query: String!, $after: String) { + orders(first: 100, after: $after, query: $query, sortKey: CREATED_AT) { + nodes { id name createdAt cancelledAt displayFinancialStatus displayFulfillmentStatus } + pageInfo { hasNextPage endCursor } + } +} +GQL; + + $orders = []; + $after = null; + do { + $data = shopify_graphql_query($env, $query, ['query' => $searchQuery, 'after' => $after]); + $connection = $data['orders'] ?? []; + foreach (($connection['nodes'] ?? []) as $order) { + $orders[] = $order; + } + $after = ($connection['pageInfo']['hasNextPage'] ?? false) ? ($connection['pageInfo']['endCursor'] ?? null) : null; + } while ($after !== null); + + return $orders; +} + +function shopify_delta_load_order(array $env, string $orderGid): array +{ + $query = <<<'GQL' +query ($id: ID!) { + order(id: $id) { + id name createdAt updatedAt cancelledAt displayFinancialStatus displayFulfillmentStatus + currentSubtotalPriceSet { shopMoney { amount currencyCode } } + totalShippingPriceSet { shopMoney { amount currencyCode } } + totalTaxSet { shopMoney { amount currencyCode } } + totalPriceSet { shopMoney { amount currencyCode } } + customer { id firstName lastName email phone } + billingAddress { firstName lastName company address1 address2 zip city provinceCode country countryCodeV2 phone } + shippingAddress { firstName lastName company address1 address2 zip city provinceCode country countryCodeV2 phone } + lineItems(first: 50) { + nodes { id sku name quantity originalUnitPriceSet { shopMoney { amount currencyCode } } variant { id sku product { id } } } + } + } +} +GQL; + $data = shopify_graphql_query($env, $query, ['id' => $orderGid]); + $order = $data['order'] ?? null; + if (!is_array($order)) { + throw new RuntimeException("Shopify-Bestellung nicht gefunden: {$orderGid}"); + } + return $order; +} + +function run_shopify_delta(array $env, PDO $pdo, string $cutoff, bool $dryRun): array +{ + $cutoffDate = new DateTimeImmutable($cutoff); + $searchQuery = 'created_at:>' . $cutoffDate->format('Y-m-d'); + $counters = [ + 'status' => 'done', + 'cutoff_timestamp' => $cutoffDate->format(DATE_ATOM), + 'dry_run' => $dryRun, + 'checked_orders' => 0, + 'skipped_cancelled_or_refunded' => 0, + 'already_imported' => 0, + 'handed_off' => 0, + 'clarification_order_gids' => [], + ]; + + foreach (shopify_delta_load_orders($env, $searchQuery) as $summaryOrder) { + $createdAt = new DateTimeImmutable((string) ($summaryOrder['createdAt'] ?? '')); + if ($createdAt <= $cutoffDate) { + continue; + } + $counters['checked_orders']++; + if (($summaryOrder['cancelledAt'] ?? null) !== null || strtoupper((string) ($summaryOrder['displayFinancialStatus'] ?? '')) === 'REFUNDED') { + $counters['skipped_cancelled_or_refunded']++; + continue; + } + + try { + $order = shopify_delta_load_order($env, (string) $summaryOrder['id']); + $result = project_shopify_order($pdo, $order, $env, $dryRun); + if (($result['status'] ?? '') === 'already_imported') { + $counters['already_imported']++; + } else { + $counters['handed_off']++; + } + } catch (Throwable $error) { + $counters['status'] = 'partial_success'; + $counters['clarification_order_gids'][] = (string) $summaryOrder['id']; + } + } + + return $counters; +} + +if (PHP_SAPI === 'cli' && realpath($_SERVER['SCRIPT_FILENAME'] ?? '') === __FILE__) { + $argv = $_SERVER['argv'] ?? []; + $cutoff = shopify_delta_argument('cutoff', $argv); + $execute = in_array('--execute', $argv, true); + $dryRun = !$execute; + if ($cutoff === null || $cutoff === '') { + fwrite(STDERR, "Usage: php shopify-delta-run.php --cutoff=ISO-8601 [--execute]\n"); + exit(2); + } + + $env = expand_env_values(parse_env_file(__DIR__ . '/../../../.env')); + $pdo = connect_database($env); + $runId = shopify_delta_start_run($pdo, $cutoff, $dryRun); + try { + $result = run_shopify_delta($env, $pdo, $cutoff, $dryRun); + $result['technical_run_id'] = $runId; + shopify_delta_finish_run($pdo, $runId, $result['status'], $result); + echo json_encode($result, JSON_UNESCAPED_UNICODE | JSON_UNESCAPED_SLASHES) . PHP_EOL; + exit($result['status'] === 'done' ? 0 : 1); + } catch (Throwable $error) { + shopify_delta_finish_run($pdo, $runId, 'failed', ['status' => 'failed', 'technical_run_id' => $runId], $error); + fwrite(STDERR, $error->getMessage() . PHP_EOL); + exit(1); + } +}