From 70e00e2630bf5636eee313558ffbe6c363495892 Mon Sep 17 00:00:00 2001 From: Sacha Al Himdani Date: Wed, 19 Nov 2025 18:15:56 +0100 Subject: [PATCH] Fix snapshots Signed-off-by: Sacha Al Himdani --- pkg/coredata/asset_vendor.go | 3 ++- pkg/coredata/datum_vendor.go | 3 ++- pkg/coredata/processing_activity_vendor.go | 17 ++++++++++++----- pkg/probo/processing_activity_service.go | 4 ++-- 4 files changed, 18 insertions(+), 9 deletions(-) diff --git a/pkg/coredata/asset_vendor.go b/pkg/coredata/asset_vendor.go index 488a60d74..c87dc5f5e 100644 --- a/pkg/coredata/asset_vendor.go +++ b/pkg/coredata/asset_vendor.go @@ -154,11 +154,12 @@ WITH FROM asset_vendors WHERE %s AND asset_id = ANY(SELECT id FROM source_assets) AND snapshot_id IS NULL ) -INSERT INTO asset_vendors (tenant_id, asset_id, vendor_id, snapshot_id, created_at) +INSERT INTO asset_vendors (tenant_id, asset_id, vendor_id, organization_id, snapshot_id, created_at) SELECT @tenant_id, sa.id, sv.id, + @organization_id, @snapshot_id, av.created_at FROM source_asset_vendors av diff --git a/pkg/coredata/datum_vendor.go b/pkg/coredata/datum_vendor.go index f96fb6bae..67b5d598d 100644 --- a/pkg/coredata/datum_vendor.go +++ b/pkg/coredata/datum_vendor.go @@ -149,11 +149,12 @@ WITH FROM data_vendors WHERE %s AND datum_id = ANY(SELECT id FROM source_data) ) -INSERT INTO data_vendors (tenant_id, datum_id, vendor_id, snapshot_id, created_at) +INSERT INTO data_vendors (tenant_id, datum_id, vendor_id, organization_id, snapshot_id, created_at) SELECT @tenant_id, sd.id, sv.id, + @organization_id, @snapshot_id, dv.created_at FROM source_data_vendors dv diff --git a/pkg/coredata/processing_activity_vendor.go b/pkg/coredata/processing_activity_vendor.go index 652392594..d766aa6d8 100644 --- a/pkg/coredata/processing_activity_vendor.go +++ b/pkg/coredata/processing_activity_vendor.go @@ -20,9 +20,9 @@ import ( "maps" "time" - "go.probo.inc/probo/pkg/gid" "github.com/jackc/pgx/v5" "go.gearno.de/kit/pg" + "go.probo.inc/probo/pkg/gid" ) type ( @@ -46,6 +46,7 @@ func (pav ProcessingActivityVendors) Merge( conn pg.Conn, scope Scoper, processingActivityID gid.GID, + organizationID gid.GID, vendorIDs []gid.GID, ) error { q := ` @@ -54,6 +55,7 @@ WITH vendor_ids AS ( unnest(@vendor_ids::text[]) AS vendor_id, @tenant_id AS tenant_id, @processing_activity_id AS processing_activity_id, + @organization_id AS organization_id, @created_at::timestamptz AS created_at ) MERGE INTO processing_activity_vendors AS tgt @@ -62,8 +64,8 @@ ON tgt.tenant_id = src.tenant_id AND tgt.processing_activity_id = src.processing_activity_id AND tgt.vendor_id = src.vendor_id WHEN NOT MATCHED - THEN INSERT (tenant_id, processing_activity_id, vendor_id, created_at) - VALUES (src.tenant_id, src.processing_activity_id, src.vendor_id, src.created_at) + THEN INSERT (tenant_id, processing_activity_id, vendor_id, organization_id, created_at) + VALUES (src.tenant_id, src.processing_activity_id, src.vendor_id, src.organization_id, src.created_at) WHEN NOT MATCHED BY SOURCE AND tgt.tenant_id = @tenant_id AND tgt.processing_activity_id = @processing_activity_id THEN DELETE @@ -72,6 +74,7 @@ WHEN NOT MATCHED args := pgx.StrictNamedArgs{ "tenant_id": scope.GetTenantID(), "processing_activity_id": processingActivityID, + "organization_id": organizationID, "created_at": time.Now(), "vendor_ids": vendorIDs, } @@ -89,17 +92,19 @@ func (pav ProcessingActivityVendors) Insert( conn pg.Conn, scope Scoper, processingActivityID gid.GID, + organizationID gid.GID, vendorIDs []gid.GID, ) error { q := ` WITH vendor_ids AS ( SELECT unnest(@vendor_ids::text[]) AS vendor_id ) -INSERT INTO processing_activity_vendors (tenant_id, processing_activity_id, vendor_id, created_at) +INSERT INTO processing_activity_vendors (tenant_id, processing_activity_id, vendor_id, organization_id, created_at) SELECT @tenant_id AS tenant_id, @processing_activity_id AS processing_activity_id, vendor_id, + @organization_id AS organization_id, @created_at AS created_at FROM vendor_ids ` @@ -107,6 +112,7 @@ FROM vendor_ids args := pgx.StrictNamedArgs{ "tenant_id": scope.GetTenantID(), "processing_activity_id": processingActivityID, + "organization_id": organizationID, "created_at": time.Now(), "vendor_ids": vendorIDs, } @@ -148,11 +154,12 @@ WITH FROM processing_activity_vendors WHERE %s AND processing_activity_id = ANY(SELECT id FROM source_processing_activities) AND snapshot_id IS NULL ) -INSERT INTO processing_activity_vendors (tenant_id, processing_activity_id, vendor_id, snapshot_id, created_at) +INSERT INTO processing_activity_vendors (tenant_id, processing_activity_id, vendor_id, organization_id, snapshot_id, created_at) SELECT @tenant_id, spa.id, sv.id, + @organization_id, @snapshot_id, pav.created_at FROM source_processing_activity_vendors pav diff --git a/pkg/probo/processing_activity_service.go b/pkg/probo/processing_activity_service.go index b4c60b7b9..139fcdbdf 100644 --- a/pkg/probo/processing_activity_service.go +++ b/pkg/probo/processing_activity_service.go @@ -185,7 +185,7 @@ func (s *ProcessingActivityService) Create( } if len(req.VendorIDs) > 0 { - if err := processingActivityVendors.Insert(ctx, conn, s.svc.scope, processingActivity.ID, req.VendorIDs); err != nil { + if err := processingActivityVendors.Insert(ctx, conn, s.svc.scope, processingActivity.ID, req.OrganizationID, req.VendorIDs); err != nil { return fmt.Errorf("cannot create processing activity vendors: %w", err) } } @@ -268,7 +268,7 @@ func (s *ProcessingActivityService) Update( } if req.VendorIDs != nil { - if err := processingActivityVendors.Merge(ctx, conn, s.svc.scope, processingActivity.ID, *req.VendorIDs); err != nil { + if err := processingActivityVendors.Merge(ctx, conn, s.svc.scope, processingActivity.ID, processingActivity.OrganizationID, *req.VendorIDs); err != nil { return fmt.Errorf("cannot update processing activity vendors: %w", err) } }