From 0782a9adb599ad28418560e1db6fd581a4403cd6 Mon Sep 17 00:00:00 2001 From: Amirali Daliri Date: Sun, 13 Sep 2026 16:51:40 +0330 Subject: [PATCH 1/2] Fix Clojure v4 strict load error handling Propagate dynamic load settings into executor workers, abort strict COPY failures after rollback, and surface MSSQL JDBC read errors instead of treating them as EOF or NULL. Add focused unit coverage for strict versus resume behavior and source read failures. --- .../workflows/clojure-integration-tests.yml | 4 +- clojure/Makefile | 4 +- clojure/src/pgloader/core.clj | 7 +- clojure/src/pgloader/prefetch.clj | 8 ++- clojure/src/pgloader/source/mssql.clj | 15 +--- clojure/test/pgloader/prefetch_test.clj | 71 +++++++++++++++++++ clojure/test/pgloader/source/mssql_test.clj | 64 +++++++++++++++++ 7 files changed, 153 insertions(+), 20 deletions(-) create mode 100644 clojure/test/pgloader/prefetch_test.clj create mode 100644 clojure/test/pgloader/source/mssql_test.clj diff --git a/.github/workflows/clojure-integration-tests.yml b/.github/workflows/clojure-integration-tests.yml index d08288d7..babf1bf3 100644 --- a/.github/workflows/clojure-integration-tests.yml +++ b/.github/workflows/clojure-integration-tests.yml @@ -83,12 +83,14 @@ jobs: clojure -M:test -m cognitect.test-runner \ -n pgloader.cast-test \ -n pgloader.batch-test \ + -n pgloader.prefetch-test \ -n pgloader.ddl-test \ -n pgloader.ddl.citus-test \ -n pgloader.load-file.parser-test \ -n pgloader.transforms-test \ -n pgloader.pg-service-test \ - -n pgloader.cli-test + -n pgloader.cli-test \ + -n pgloader.source.mssql-test # ── Build test-runner Docker image (contains v3 pgloader binary) ─────────── # Built once and shared with all integration jobs via artifact. diff --git a/clojure/Makefile b/clojure/Makefile index 9c917c58..1d5185a3 100644 --- a/clojure/Makefile +++ b/clojure/Makefile @@ -48,12 +48,14 @@ test-unit: clojure -M:test -m cognitect.test-runner \ -n pgloader.cast-test \ -n pgloader.batch-test \ + -n pgloader.prefetch-test \ -n pgloader.ddl-test \ -n pgloader.ddl.citus-test \ -n pgloader.load-file.parser-test \ -n pgloader.transforms-test \ -n pgloader.pg-service-test \ - -n pgloader.cli-test + -n pgloader.cli-test \ + -n pgloader.source.mssql-test # ─── E2E integration tests ──────────────────────────────────────────────────── # All suite management lives in tests/Makefile. diff --git a/clojure/src/pgloader/core.clj b/clojure/src/pgloader/core.clj index 45518433..abee5685 100644 --- a/clojure/src/pgloader/core.clj +++ b/clojure/src/pgloader/core.clj @@ -786,7 +786,7 @@ (fn [[_i t]] (.submit ^ExecutorService workers-pool ^java.util.concurrent.Callable - (fn [] + (bound-fn [] (let [worker-src (source-from-uri source-uri table-spec with-options source-overrides (:decoding-as cmd)) worker-pg (postgres-connection target-uri) @@ -806,7 +806,8 @@ table (:table-name t) cols (:columns t) table-label (table-stats-label schema table)] - (when-not (@failed-tables table) + (when-not (or (@failed-tables table) + (and copy/*on-error-stop* @load-failed)) (let [;; Generated columns exist on the target (DDL emits ;; GENERATED ALWAYS AS) but must be excluded from COPY ;; since PostgreSQL cannot accept values for them. @@ -835,7 +836,7 @@ (mapv (fn [part-src] (.submit ^ExecutorService part-exec ^java.util.concurrent.Callable - (fn [] + (bound-fn [] (let [part-pg (postgres-connection target-uri)] (when pg-params (doseq [param pg-params] diff --git a/clojure/src/pgloader/prefetch.clj b/clojure/src/pgloader/prefetch.clj index c973a292..2af36bd0 100644 --- a/clojure/src/pgloader/prefetch.clj +++ b/clojure/src/pgloader/prefetch.clj @@ -75,7 +75,7 @@ (defn- send-batch-or-retry! "Send a single batch, or handle errors and return updated counters. Returns {:status :ok :rows-ok ... :errors ... :ws-nanos ... :bytes ... :reject-paths ...} - for success/retry, or throws for non-retryable errors." + for success/retry. In strict mode, rolls back and propagates COPY errors." [^PGConnection pg-conn table-spec ^String copy-sql-str b rows-ok errors ws-nanos bytes reject-paths] (let [batch-start (System/nanoTime)] @@ -96,6 +96,8 @@ (throw e)) (catch PSQLException e (.rollback ^Connection pg-conn) + (when copy/*on-error-stop* + (throw e)) (log/info "Entering error recovery.") (let [retry-result (batch/retry-batch! b table-spec e pg-conn)] {:status :retry @@ -112,8 +114,8 @@ (defn writer-task "Virtual thread task that drains batches from the pipeline queue and sends them to PostgreSQL via CopyManager. - Each batch gets its own transaction. On data errors, retry-batch! - handles per-row recovery with independent sub-batch commits. + Each batch gets its own transaction. In resume mode, retry-batch! handles + per-row recovery with independent sub-batch commits. Returns {:rows-ok n :rows-bad n :ws-nanos n :bytes n :reject-paths {...}}." [^PGConnection pg-conn table-spec ^CopyPipeline pipeline] (let [copy-sql-str (copy/copy-sql table-spec) diff --git a/clojure/src/pgloader/source/mssql.clj b/clojure/src/pgloader/source/mssql.clj index 87661007..c5fc123a 100644 --- a/clojure/src/pgloader/source/mssql.clj +++ b/clojure/src/pgloader/source/mssql.clj @@ -215,24 +215,15 @@ meta (.getMetaData rs) n (.getColumnCount meta)] ((fn thisfn [] - (when (try (.next rs) - (catch Exception e - (log/warn (str "MSSQL row advance error in " table-name ": " (.getMessage e))) - false)) + (when (.next rs) (lazy-seq (cons (loop [i 1 result (transient [])] (if (<= i n) (recur (inc i) (conj! result ;; getString preserves full decimal/numeric precision (#1615, #1619). - ;; Per-column error recovery: substitute nil on any error. - (try - (.getString rs i) - (catch Exception e - (log/warn (str "MSSQL column " i " read error in " - table-name " (substituting NULL): " - (.getMessage e))) - nil)))) + ;; Driver errors must abort rather than change data to NULL. + (.getString rs i))) (persistent! result))) (thisfn))))))) (catch Exception e diff --git a/clojure/test/pgloader/prefetch_test.clj b/clojure/test/pgloader/prefetch_test.clj new file mode 100644 index 00000000..39c15d7d --- /dev/null +++ b/clojure/test/pgloader/prefetch_test.clj @@ -0,0 +1,71 @@ +(ns pgloader.prefetch-test + (:require [clojure.test :refer [deftest is testing]] + [pgloader.batch :as batch] + [pgloader.copy :as copy] + [pgloader.prefetch :as prefetch]) + (:import [java.lang.reflect InvocationHandler Proxy] + [java.nio.charset StandardCharsets] + [java.sql Connection] + [org.postgresql PGConnection] + [org.postgresql.util PSQLException PSQLState])) + +(defn- pg-connection + [rolled-back?] + (Proxy/newProxyInstance + (.getClassLoader PGConnection) + (into-array Class [PGConnection Connection]) + (reify InvocationHandler + (invoke [_ _ method _] + (when (= "rollback" (.getName method)) + (reset! rolled-back? true)))))) + +(deftest on-error-stop-does-not-retry-a-failed-copy-batch + (let [rolled-back? (atom false) + retried? (atom false) + test-batch (batch/batch-add-row! + (batch/make-batch 10 1024) + (.getBytes "1\n" StandardCharsets/UTF_8)) + error (PSQLException. "bad value" PSQLState/DATA_ERROR)] + (with-redefs [batch/send-batch! (fn [& _] (throw error)) + batch/retry-batch! (fn [& _] + (reset! retried? true) + {:rows-ok 0 :errors 1})] + (binding [copy/*on-error-stop* true] + (is (thrown-with-msg? + PSQLException #"bad value" + (#'prefetch/send-batch-or-retry! + (pg-connection rolled-back?) + {:target-schema "public" :target-table "items"} + "COPY public.items FROM STDIN" + test-batch 0 0 0 0 nil))))) + (testing "the active transaction is rolled back before propagation" + (is @rolled-back?)) + (testing "strict mode never enters row rejection" + (is (false? @retried?))))) + +(deftest resume-mode-still-retries-a-failed-copy-batch + (let [rolled-back? (atom false) + retried? (atom false) + test-batch (batch/batch-add-row! + (batch/make-batch 10 1024) + (.getBytes "1\n" StandardCharsets/UTF_8)) + error (PSQLException. "bad value" PSQLState/DATA_ERROR)] + (with-redefs [batch/send-batch! (fn [& _] (throw error)) + batch/retry-batch! (fn [& _] + (reset! retried? true) + {:rows-ok 1 :errors 1})] + (binding [copy/*on-error-stop* false] + (is (= {:status :retry + :rows-ok 1 + :errors 1 + :bytes 2 + :reject-paths nil} + (select-keys + (#'prefetch/send-batch-or-retry! + (pg-connection rolled-back?) + {:target-schema "public" :target-table "items"} + "COPY public.items FROM STDIN" + test-batch 0 0 0 0 nil) + [:status :rows-ok :errors :bytes :reject-paths]))))) + (is @rolled-back?) + (is @retried?))) diff --git a/clojure/test/pgloader/source/mssql_test.clj b/clojure/test/pgloader/source/mssql_test.clj new file mode 100644 index 00000000..e1a35023 --- /dev/null +++ b/clojure/test/pgloader/source/mssql_test.clj @@ -0,0 +1,64 @@ +(ns pgloader.source.mssql-test + (:require [clojure.test :refer [deftest is]] + [pgloader.source.mssql :as mssql] + [pgloader.source.protocol :as source]) + (:import [java.lang.reflect InvocationHandler Proxy] + [java.sql Connection PreparedStatement ResultSet + ResultSetMetaData SQLException])) + +(defn- interface-proxy + [interface invoke-fn] + (Proxy/newProxyInstance + (.getClassLoader interface) + (into-array Class [interface]) + (reify InvocationHandler + (invoke [_ _ method args] + (invoke-fn (.getName method) args))))) + +(defn- failing-source + [failure-point] + (let [advance-count (atom 0) + metadata (interface-proxy + ResultSetMetaData + (fn [method _] + (case method + "getColumnCount" (int 1) + nil))) + result-set (interface-proxy + ResultSet + (fn [method _] + (case method + "next" (if (= failure-point :row-advance) + (throw (SQLException. "row advance failed")) + (= 1 (swap! advance-count inc))) + "getMetaData" metadata + "getString" (throw (SQLException. "column read failed")) + nil))) + statement (interface-proxy + PreparedStatement + (fn [method _] + (case method + "executeQuery" result-set + nil))) + connection (interface-proxy + Connection + (fn [method _] + (case method + "prepareStatement" statement + nil)))] + (mssql/->MSSQLSource connection "fixture"))) + +(def table-spec + {:table-name "items" + :schema "dbo" + :columns [{:column-name "value"}]}) + +(deftest row-advance-errors-propagate + (is (thrown-with-msg? + SQLException #"row advance failed" + (doall (source/read-rows (failing-source :row-advance) table-spec))))) + +(deftest column-read-errors-propagate + (is (thrown-with-msg? + SQLException #"column read failed" + (doall (source/read-rows (failing-source :column-read) table-spec))))) From 113c276d1e88343bf19b4756d9c66ba9a49c220f Mon Sep 17 00:00:00 2001 From: Amirali Daliri Date: Sun, 13 Sep 2026 17:14:19 +0330 Subject: [PATCH 2/2] Harden strict worker failure cleanup --- .../workflows/clojure-integration-tests.yml | 1 + clojure/Makefile | 1 + clojure/src/pgloader/core.clj | 25 ++++++- clojure/src/pgloader/prefetch.clj | 31 ++++++--- clojure/test/pgloader/cli_test.clj | 19 ++++- clojure/test/pgloader/core_test.clj | 24 +++++++ clojure/test/pgloader/prefetch_test.clj | 69 ++++++++++++++++--- 7 files changed, 147 insertions(+), 23 deletions(-) create mode 100644 clojure/test/pgloader/core_test.clj diff --git a/.github/workflows/clojure-integration-tests.yml b/.github/workflows/clojure-integration-tests.yml index babf1bf3..4a631427 100644 --- a/.github/workflows/clojure-integration-tests.yml +++ b/.github/workflows/clojure-integration-tests.yml @@ -90,6 +90,7 @@ jobs: -n pgloader.transforms-test \ -n pgloader.pg-service-test \ -n pgloader.cli-test \ + -n pgloader.core-test \ -n pgloader.source.mssql-test # ── Build test-runner Docker image (contains v3 pgloader binary) ─────────── diff --git a/clojure/Makefile b/clojure/Makefile index 1d5185a3..241c8a96 100644 --- a/clojure/Makefile +++ b/clojure/Makefile @@ -55,6 +55,7 @@ test-unit: -n pgloader.transforms-test \ -n pgloader.pg-service-test \ -n pgloader.cli-test \ + -n pgloader.core-test \ -n pgloader.source.mssql-test # ─── E2E integration tests ──────────────────────────────────────────────────── diff --git a/clojure/src/pgloader/core.clj b/clojure/src/pgloader/core.clj index abee5685..4faaac7b 100644 --- a/clojure/src/pgloader/core.clj +++ b/clojure/src/pgloader/core.clj @@ -27,7 +27,8 @@ (:import [org.postgresql PGConnection] [java.sql Connection DriverManager] [java.io File] - [java.util.concurrent Executors ExecutorService Future TimeUnit]) + [java.util.concurrent Executors ExecutorService Future TimeUnit + ExecutionException]) (:require [clojure.tools.logging :as log])) (set! *warn-on-reflection* true) @@ -408,6 +409,25 @@ :reset-sequences true} with-options)) +(defn- await-table-futures! + "Wait for every table worker, finish index work, then propagate a strict failure." + [table-futs ^ExecutorService idx-executor] + (let [first-failure (volatile! nil)] + (doseq [^Future f table-futs] + (try + (.get f) + (catch ExecutionException e + (when-not @first-failure + (vreset! first-failure (or (.getCause e) e)))) + (catch Exception e + (when-not @first-failure + (vreset! first-failure e))))) + (when idx-executor + (.shutdown idx-executor) + (.awaitTermination idx-executor Long/MAX_VALUE TimeUnit/NANOSECONDS)) + (when (and copy/*on-error-stop* @first-failure) + (throw @first-failure)))) + (defn run-command [cmd opts] (if (= :archive (:load-type cmd)) @@ -950,8 +970,7 @@ (map-indexed vector cat))] (.shutdown ^ExecutorService workers-pool) (.awaitTermination ^ExecutorService workers-pool Long/MAX_VALUE TimeUnit/NANOSECONDS) - (doseq [^Future f table-futs] - (try (.get f) (catch Exception _))) + (await-table-futures! table-futs idx-executor) (stats/update-entry! :post "COPY Wall-Clock Time" :rows workers :bytes (:bytes (stats/get-totals :data)) diff --git a/clojure/src/pgloader/prefetch.clj b/clojure/src/pgloader/prefetch.clj index 2af36bd0..936ff44e 100644 --- a/clojure/src/pgloader/prefetch.clj +++ b/clojure/src/pgloader/prefetch.clj @@ -4,7 +4,7 @@ [org.postgresql.util PSQLException] [java.sql Connection] [java.util.concurrent LinkedBlockingQueue - BlockingQueue] + BlockingQueue TimeUnit] [java.util.concurrent.atomic AtomicBoolean AtomicLong] [java.nio.charset StandardCharsets]) (:require [pgloader.batch :as batch] @@ -95,7 +95,11 @@ (if cause (.getMessage ^Throwable cause) "unknown")))) (throw e)) (catch PSQLException e - (.rollback ^Connection pg-conn) + (try + (.rollback ^Connection pg-conn) + (catch Exception rollback-error + (.addSuppressed e rollback-error) + (throw e))) (when copy/*on-error-stop* (throw e)) (log/info "Entering error recovery.") @@ -125,20 +129,25 @@ ws-nanos (long 0) bytes (long 0) reject-paths nil] - (let [item (.take ^BlockingQueue (.queue pipeline))] - (if (= :end-of-data item) + (let [item (.poll ^BlockingQueue (.queue pipeline) + 100 TimeUnit/MILLISECONDS)] + (if (or (= :end-of-data item) + (and (nil? item) + (.get ^AtomicBoolean (.done pipeline)))) {:rows-ok rows-ok :rows-bad errors :ws-nanos (- (System/nanoTime) start) :bytes bytes :reject-paths reject-paths} - (let [^batch/Batch b item - result (send-batch-or-retry! - pg-conn table-spec copy-sql-str - b rows-ok errors ws-nanos bytes reject-paths)] - (recur (long (:rows-ok result)) (long (:errors result)) - (long (:ws-nanos result)) (long (:bytes result)) - (:reject-paths result)))))))) + (if (nil? item) + (recur rows-ok errors ws-nanos bytes reject-paths) + (let [^batch/Batch b item + result (send-batch-or-retry! + pg-conn table-spec copy-sql-str + b rows-ok errors ws-nanos bytes reject-paths)] + (recur (long (:rows-ok result)) (long (:errors result)) + (long (:ws-nanos result)) (long (:bytes result)) + (:reject-paths result))))))))) (defn copy-table! "Orchestrate the full copy of a single table. diff --git a/clojure/test/pgloader/cli_test.clj b/clojure/test/pgloader/cli_test.clj index 72c7ab47..c47075f2 100644 --- a/clojure/test/pgloader/cli_test.clj +++ b/clojure/test/pgloader/cli_test.clj @@ -1,6 +1,8 @@ (ns pgloader.cli-test (:require [clojure.test :refer [deftest is testing]] - [pgloader.cli :as cli])) + [pgloader.cli :as cli] + [pgloader.core :as core] + [pgloader.load-file.parser :as parser])) (deftest test-parse-args-basic (testing "positional .load file" @@ -81,3 +83,18 @@ (is (= ["quote identifiers" "include drop"] (:with-opts opts))) (is (= "sqlite:///tmp/test.db" (:source-uri opts))) (is (= "pgsql:///target" (:target-uri opts)))))) + +(deftest strict-load-file-failure-stops-later-files + (let [calls (atom []) + failure (ex-info "copy failed" {:file "first.load"})] + (with-redefs [parser/parse-file (fn [file] {:ok {:file file}}) + core/run-command (fn [cmd _opts] + (swap! calls conj (:file cmd)) + (when (= "first.load" (:file cmd)) + (throw failure)))] + (let [thrown (try + (cli/run ["first.load" "second.load"]) + nil + (catch Exception e e))] + (is (identical? failure thrown)) + (is (= ["first.load"] @calls)))))) diff --git a/clojure/test/pgloader/core_test.clj b/clojure/test/pgloader/core_test.clj new file mode 100644 index 00000000..7a1e609e --- /dev/null +++ b/clojure/test/pgloader/core_test.clj @@ -0,0 +1,24 @@ +(ns pgloader.core-test + (:require [clojure.test :refer [deftest is testing]] + [pgloader.copy :as copy] + [pgloader.core :as core]) + (:import [java.util.concurrent CompletableFuture Executors ExecutorService])) + +(deftest await-table-futures-propagates-the-first-strict-failure + (let [failure (ex-info "copy failed" {:table "items"}) + failed (CompletableFuture/failedFuture failure) + complete (CompletableFuture/completedFuture :ok)] + (testing "strict mode unwraps and propagates a worker failure" + (let [^ExecutorService index-executor (Executors/newSingleThreadExecutor)] + (.submit index-executor ^Runnable (fn [] nil)) + (binding [copy/*on-error-stop* true] + (let [thrown (try + (#'core/await-table-futures! [failed complete] index-executor) + nil + (catch Exception e e))] + (is (identical? failure thrown)) + (is (.isShutdown index-executor)) + (is (.isTerminated index-executor)))))) + (testing "resume mode still waits for failures without aborting the load" + (binding [copy/*on-error-stop* false] + (is (nil? (#'core/await-table-futures! [failed complete] nil))))))) diff --git a/clojure/test/pgloader/prefetch_test.clj b/clojure/test/pgloader/prefetch_test.clj index 39c15d7d..9358fc2f 100644 --- a/clojure/test/pgloader/prefetch_test.clj +++ b/clojure/test/pgloader/prefetch_test.clj @@ -6,18 +6,45 @@ (:import [java.lang.reflect InvocationHandler Proxy] [java.nio.charset StandardCharsets] [java.sql Connection] + [java.util.concurrent.atomic AtomicBoolean] [org.postgresql PGConnection] [org.postgresql.util PSQLException PSQLState])) (defn- pg-connection - [rolled-back?] - (Proxy/newProxyInstance - (.getClassLoader PGConnection) - (into-array Class [PGConnection Connection]) - (reify InvocationHandler - (invoke [_ _ method _] - (when (= "rollback" (.getName method)) - (reset! rolled-back? true)))))) + ([rolled-back?] (pg-connection rolled-back? nil)) + ([rolled-back? rollback-error] + (Proxy/newProxyInstance + (.getClassLoader PGConnection) + (into-array Class [PGConnection Connection]) + (reify InvocationHandler + (invoke [_ _ method _] + (when (= "rollback" (.getName method)) + (reset! rolled-back? true) + (when rollback-error + (throw rollback-error)))))))) + +(deftest writer-finishes-when-reader-fails-with-a-full-queue + (binding [copy/*prefetch-queue-capacity* 1] + (let [pipeline (prefetch/make-pipeline 10 1024) + test-batch (batch/batch-add-row! + (batch/make-batch 10 1024) + (.getBytes "1\n" StandardCharsets/UTF_8)) + writer (atom nil)] + (.put (.queue pipeline) test-batch) + (is (zero? (.remainingCapacity (.queue pipeline)))) + (.set ^AtomicBoolean (.done pipeline) true) + (with-redefs [copy/copy-sql (constantly "COPY public.items FROM STDIN") + batch/send-batch! (fn [& _] {:rows 1})] + (reset! writer + (future + (prefetch/writer-task + (pg-connection (atom false)) + {:target-schema "public" :target-table "items"} + pipeline))) + (try + (is (= 1 (:rows-ok (deref @writer 500 ::timed-out)))) + (finally + (future-cancel @writer))))))) (deftest on-error-stop-does-not-retry-a-failed-copy-batch (let [rolled-back? (atom false) @@ -69,3 +96,29 @@ [:status :rows-ok :errors :bytes :reject-paths]))))) (is @rolled-back?) (is @retried?))) + +(deftest rollback-failure-preserves-the-copy-error-and-aborts + (let [rolled-back? (atom false) + retried? (atom false) + test-batch (batch/batch-add-row! + (batch/make-batch 10 1024) + (.getBytes "1\n" StandardCharsets/UTF_8)) + copy-error (PSQLException. "bad value" PSQLState/DATA_ERROR) + rollback-error (java.sql.SQLException. "rollback failed")] + (with-redefs [batch/send-batch! (fn [& _] (throw copy-error)) + batch/retry-batch! (fn [& _] + (reset! retried? true) + {:rows-ok 0 :errors 1})] + (binding [copy/*on-error-stop* false] + (let [thrown (try + (#'prefetch/send-batch-or-retry! + (pg-connection rolled-back? rollback-error) + {:target-schema "public" :target-table "items"} + "COPY public.items FROM STDIN" + test-batch 0 0 0 0 nil) + nil + (catch PSQLException e e))] + (is (identical? copy-error thrown)) + (is (= [rollback-error] (vec (.getSuppressed thrown))))))) + (is @rolled-back?) + (is (false? @retried?))))