diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..7cd94ef --- /dev/null +++ b/.gitignore @@ -0,0 +1,13 @@ +/target +/classes +/checkouts +profiles.clj +pom.xml +pom.xml.asc +*.jar +*.class +/.lein-* +/.nrepl-port +.hgignore +.hg/ +.cpcache/ An exception is the =jsonb= type, +because the binary format requires a version signifier. Wrapping a +JSON string in a =JsonB= handles that, which is provided by the +library. + +** Arrays + +Impemented for the following JVM-typed arrays for: + +| JVM type | Postgres type | +|------------------+-----------------------------------| +| int[] | int4[] (aka integer[]) | +| long[] | int8[] (aka bigint[]) | +| float[] | float4[] (aka real[]) | +| double[] | float8[] (aka double precision[]) | +| byte[] | bytea | +| String[] | text[] (or varchar) | +| java.util.UUID[] | uuid[] | + + +Currently, only 1-dimensional Postgres arrays are supported. + +** TODO + +- hstore (wrapper?) +- inet, cidr, macaddr, macaddr8 +- bit strings +- composite types +- range types +- more array types? (date, timestamp, etc) diff --git a/deps.edn b/deps.edn new file mode 100644 index 0000000..16415cd --- /dev/null +++ b/deps.edn @@ -0,0 +1,12 @@ +{:paths ["src"] + :deps {org.postgresql/postgresql {:mvn/version "42.2.6"}} + :aliases {:run-tests + {:extra-paths ["test"] + :extra-deps {lambdaisland/kaocha {:mvn/version "0.0-529"} + org.clojure/java.jdbc {:mvn/version "0.7.9"}} + :main-opts ["-m" "kaocha.runner"]} + :dev + {:extra-paths ["dev" "test"] + :extra-deps {org.clojure/tools.namespace {:mvn/version "0.3.0"} + org.clojure/java.jdbc {:mvn/version "0.7.9"}} + :jvm-opts ["-XX:-OmitStackTraceInFastThrow"]}}} diff --git a/dev/user.clj b/dev/user.clj new file mode 100644 index 0000000..4d153a7 --- /dev/null +++ b/dev/user.clj @@ -0,0 +1,19 @@ +(ns ^{:clojure.tools.namespace.repl/load false + :clojure.tools.namespace.repl/unload false} user + (:require + ;; Defaults copied from clojure.main + [clojure.repl :refer (source apropos dir pst doc find-doc)] + [clojure.java.javadoc :refer (javadoc)] + [clojure.pprint :refer (pp pprint)] + ;; very common convenience ns aliases + [clojure.java.io :as io] + [clojure.set :as set] + [clojure.edn :as edn] + [clojure.string :as str] + [clojure.test :as test] + [clojure.tools.namespace.repl :as ctn])) + +(ctn/set-refresh-dirs "src" "test") + +(defn reset [] + (ctn/refresh)) diff --git a/src/clj_pgcopy/core.clj b/src/clj_pgcopy/core.clj new file mode 100644 index 0000000..b169cae --- /dev/null +++ b/src/clj_pgcopy/core.clj @@ -0,0 +1,423 @@ +(ns clj-pgcopy.core + (:require [clojure.string :as str] + [clj-pgcopy.time :as ptime]) + (:import (java.io ByteArrayInputStream + ByteArrayOutputStream + BufferedReader + BufferedOutputStream + DataOutputStream) + (java.nio ByteBuffer) + (java.time LocalDateTime + LocalDate + ZonedDateTime + OffsetDateTime + Instant) + (java.util UUID) + (org.postgresql.geometric PGbox + PGcircle + PGline + PGpath + PGpolygon + PGpoint) + (org.postgresql.util PGInterval) + (org.postgresql.copy CopyManager + PGCopyOutputStream) + (org.postgresql.core BaseConnection))) + +(set! *warn-on-reflection* true) + +(defprotocol IPGBinaryWrite + (pg-type [this]) + (write-to [this ^DataOutputStream out])) + +(def oids + {:bytea 17 + :text 25 + :int4 23 + :int8 20 + :int2 21 + :char 18 + :boolean 16 + :jsonb 114 + :xml 115 + :point 600 + :line 628 + :path 602 + :box 603 + :polygon 604 + :circle 705 + :float4 700 + :float8 701 + :unknown 705 + :varchar 1043 + :date 1082 + :timestamp 1114 + :timestamptz 1184 + :interval 1186 + :numeric 1700 + :uuid 2950}) + +(defn array->bytes + "Returns binary-encoded byte array when type of array can be determined, or nil" + [pg-type coll] + (let [baos ^ByteArrayOutputStream (ByteArrayOutputStream. 1024) + [_ oid] (find oids pg-type) + coll (seq coll)] + (when oid + (with-open [out ^DataOutputStream (DataOutputStream. baos)] + (.writeInt out 1) ;; dimensions (only 1 dimensional) + (.writeInt out 1) ;; nullable values allowed + (.writeInt out oid) ;; oid of collection elements + (.writeInt out (count coll)) ;; size + (.writeInt out 1) ;; use PG default + (doseq [el coll] + (write-to el out))) + (.toByteArray baos)))) + +;; Inspired by Java impl here: +;; https://github.com/bytefish/PgBulkInsert/blob/master/PgBulkInsert/src/main/java/de/bytefish/pgbulkinsert/pgsql/handlers/BigDecimalValueHandler.java +(defn numeric-components [^BigDecimal bd] + (let [unscaled ^BigInteger (.unscaledValue bd) + sign (.signum bd) + unscaled (if (= -1 sign) (.negate unscaled) unscaled) + fraction-digits ^int (.scale bd) + fraction-groups (unchecked-divide-int (unchecked-add-int fraction-digits (int 3)) + 4) + scale-remainder (mod fraction-digits 4) + [unscaled digits] (if (zero? scale-remainder) + [unscaled (list)] + ;; scale the first value + (let [result (.divideAndRemainder unscaled (.pow (BigInteger. "10") scale-remainder)) + digit (unchecked-multiply-int ^int (.intValue ^BigInteger (aget result 1)) + (int (Math/pow 10 (- 4 scale-remainder))))] + [(aget result 0) (list digit)])) + digits (loop [^BigInteger unscaled unscaled + digits digits] + (if (.equals unscaled BigInteger/ZERO) + digits + (let [result (.divideAndRemainder unscaled (BigInteger. "10000"))] + (recur + (aget result 0) + (cons (.intValue ^BigInteger (aget result 1)) digits)))))] + {:sign sign + :fraction-groups fraction-groups + :fraction-digits fraction-digits + :digits digits})) + +(extend-protocol IPGBinaryWrite + (Class/forName "[B") + (pg-type [_] :bytea) + (write-to [ba ^DataOutputStream out] + (.writeInt out (count ba)) + (doseq [b ba] + (.writeByte out b))) + + String + (pg-type [_] :text) + (write-to [string ^DataOutputStream out] + (write-to ^bytes (.getBytes string "UTF-8") out)) + + Short + (pg-type [_] :int2) + (write-to [sh ^DataOutputStream out] + (.writeInt out 2) + (.writeShort out sh)) + + Integer + (pg-type [_] :int4) + (write-to [integer ^DataOutputStream out] + (.writeInt out 4) + (.writeInt out integer)) + + Long + (pg-type [_] :int8) + (write-to [num ^DataOutputStream out] + (.writeInt out 8) + (.writeLong out num)) + + Float + (pg-type [_] :float4) + (write-to [f ^DataOutputStream out] + (.writeInt out 4) + (.writeFloat out (.floatValue f))) + + Double + (pg-type [_] :float8) + (write-to [d ^DataOutputStream out] + (.writeInt out 8) + (.writeDouble out (.doubleValue d))) + + BigDecimal + (pg-type [_] :numeric) + (write-to [bd ^DataOutputStream out] + (let [{:keys [fraction-digits fraction-groups sign digits]} + (numeric-components bd) + n-digits (int (count digits))] + (.writeInt out (int (+ 8 (* 2 n-digits)))) + (.writeShort out n-digits) + (.writeShort out (- n-digits fraction-groups 1)) + (.writeShort out (int (if (= sign 1) 0x0000 0x4000))) + (.writeShort out fraction-digits) + (doseq [digit digits] + (.writeShort out (int digit))))) + + Boolean + (pg-type [_] :boolean) + (write-to [bool ^DataOutputStream out] + (.writeInt out 1) + (.writeByte out (if bool 1 0))) + + PGInterval + (pg-type [_] :interval) + (write-to [interval ^DataOutputStream out] + (.writeInt out 16) + (let [secs (.getSeconds interval) + mins (double (.getMinutes interval)) + hours (double (.getHours interval)) + seconds (+ secs (* 60 mins) (* 60 60 hours)) + days (.getDays interval) + months (.getMonths interval) + years (.getYears interval)] + (.writeLong out (long (* 1000000 seconds))) + (.writeInt out days) + (.writeInt out (+ months (* 12 years))))) + + java.util.Date + (pg-type [_] :timestamp) + (write-to [date ^DataOutputStream out] + (write-to (ptime/java-epoch->postgres-epoch (.getTime date)) out)) + + Instant + (pg-type [_] :timestamp) + (write-to [instant ^DataOutputStream out] + (write-to (ptime/java-epoch->postgres-epoch (.toEpochMilli instant)) out)) + + LocalDate + (pg-type [_] :date) + (write-to [date ^DataOutputStream out] + (let [days (-> date + .atStartOfDay + ptime/date-time->epoch-milli + ptime/java-epoch->postgres-days + int)] + (.writeInt out 4) + (.writeInt out days))) + + java.sql.Date + (pg-type [_] :date) + (write-to [date ^DataOutputStream out] + (write-to ^LocalDate (.toLocalDate date) out)) + + java.sql.Timestamp + (pg-type [_] :timestamp) + (write-to [ts ^DataOutputStream out] + (write-to ^Instant (.toInstant ts) out)) + + ;; Do not assume UTC, needs a time zone or offset + ;; LocalDateTime + ;; (pg-type [_] :timestamp) + ;; (write-to [ldt ^DataOutputStream out] + ;; (write-to ^Instant (.toInstant ldt) out)) + + ZonedDateTime + (pg-type [_] :timestamp) + (write-to [zdt ^DataOutputStream out] + (write-to ^Instant (.toInstant zdt) out)) + + OffsetDateTime + (pg-type [_] :timestamp) + (write-to [odt ^DataOutputStream out] + (write-to ^Instant (.toInstant odt) out)) + + PGpoint + (pg-type [_] :point) + (write-to [p ^DataOutputStream out] + (.writeInt out 16) + (.writeDouble out (.-x p)) + (.writeDouble out (.-y p))) + + PGline + (pg-type [_] :line) + (write-to [line ^DataOutputStream out] + (.writeInt out 24) + (.writeDouble out (.-a line)) + (.writeDouble out (.-b line)) + (.writeDouble out (.-c line))) + + PGpath + (pg-type [_] :path) + (write-to [p ^DataOutputStream out] + (let [points (.-points p) + closed (byte (if (.-open p) 0 1)) + byte-count (+ 1 ;; open/closed + 4 ;; number of points + (* 16 (count points)) ;; point data + )] + (.writeInt out byte-count) + (.writeByte out closed) + (.writeInt out (count points)) + (doseq [^PGpoint point points] + (.writeDouble out (.-x point)) + (.writeDouble out (.-y point))))) + + PGpolygon + (pg-type [_] :polygon) + (write-to [p ^DataOutputStream out] + (let [points (.-points p) + byte-count (+ 4 (* 16 (count points)))] + (.writeInt out byte-count) + (.writeInt out (count points)) + (doseq [^PGpoint point points] + (.writeDouble out (.-x point)) + (.writeDouble out (.-y point))))) + + PGbox + (pg-type [_] :box) + (write-to [box ^DataOutputStream out] + (let [points (.-point box)] + (.writeInt out 32) + (doseq [^PGpoint point points] + (.writeDouble out (.-x point)) + (.writeDouble out (.-y point))))) + + PGcircle + (pg-type [_] :circle) + (write-to [c ^DataOutputStream out] + (let [center ^PGpoint (.-center c)] + (.writeInt out 24) + (.writeDouble out (.-x center)) + (.writeDouble out (.-y center)) + (.writeDouble out (.-radius c)))) + + ;; TODO Is this a good idea? + ;; clojure.lang.APersistentVector + ;; (pg-type [_] nil) ;; array's use sub oid + ;; (write-to [this ^DataOutputStream out] + ;; (if (and (seq this) (satisfies? IPGBinaryWrite (first this))) + ;; (if-some [ba ^bytes (array->bytes (pg-type (first this)) this)] + ;; (do + ;; (.writeInt out (count ba)) + ;; (.write out ba)) + ;; (write-to nil out)) + ;; (write-to nil out))) + + UUID + (pg-type [_] :uuid) + (write-to [uuid ^DataOutputStream out] + (let [bb ^ByteBuffer (ByteBuffer/wrap (byte-array 16))] + (.putLong bb (.getMostSignificantBits uuid)) + (.putLong bb (.getLeastSignificantBits uuid)) + (.writeInt out 16) + (.writeInt out (.getInt bb 0)) + (.writeShort out (.getShort bb 4)) + (.writeShort out (.getShort bb 6)) + (.writeLong out (.getLong bb 8)))) + + nil + (oid [_] :unknown) + (write-to [_ ^DataOutputStream out] + (.writeInt out -1))) + +;; These primitive array impls have to be added separate or the +;; compiler complains + +(defmacro extend-primitive-array + [klass pg-type] + `(extend-type ~klass + IPGBinaryWrite + (pg-type [~'_] nil) + (write-to [this# ^DataOutputStream out#] + (if-some [ba# ^{:tag ~'bytes} (array->bytes ~pg-type this#)] + (do + (.writeInt out# (count ba#)) + (.write out# ba#)) + (write-to nil out#))))) + + ;; int array +(extend-primitive-array (Class/forName "[I") :int4) + + ;; long array +(extend-primitive-array (Class/forName "[J") :int8) + +;; float array +(extend-primitive-array (Class/forName "[F") :float4) + +;; double array +(extend-primitive-array (Class/forName "[D") :float8) + +;; String array +(extend-primitive-array (Class/forName "[Ljava.lang.String;") :text) + +(extend-primitive-array (Class/forName "[Ljava.util.UUID;") :uuid) + +(deftype JsonB [^String value] + Object + (toString [_] value) + IPGBinaryWrite + (pg-type [_] :jsonb) + (write-to [_ o] + (let [out ^DataOutputStream o] + (if (str/blank? value) + (write-to out nil) + (let [ba ^bytes (.getBytes value "UTF-8")] + (.writeInt out (+ 1 (count ba))) + (.writeByte out 1) ;; jsonb protocol version + (.write out ba)))))) + +(defn copy-to-stream [^ByteArrayOutputStream stream tuples] + (with-open [out ^DataOutputStream (DataOutputStream. (BufferedOutputStream. stream))] + ;; constant header + (.writeBytes out "PGCOPY\n\377\r\n\0") + ;; header flags (no OIDs) + (.writeInt out 0) + ;; header extension (unused) + (.writeInt out 0) + (doseq [tuple tuples] + (.writeShort out (count tuple)) + (doseq [field tuple] + (write-to field out))) + ;; footer constant + (.writeShort out (short -1)) + (.flush out))) + +(defn values->copy-rows-binary ^bytes [values] + (with-open [bout ^ByteArrayOutputStream (ByteArrayOutputStream. 4096)] + (copy-to-stream bout values) + (.toByteArray bout))) + +(defn ^CopyManager copy-manager [^java.sql.Connection conn] + (let [conn (if (.isWrapperFor conn BaseConnection) + (.unwrap conn BaseConnection) + conn)] + (CopyManager. conn))) + +(defn copy-table + ([conn table] + (copy-table conn table {:csv? true :headers? false})) + ([^java.sql.Connection conn table {:keys [csv? headers?]}] + (let [manager (copy-manager conn) + copy-query (str "COPY " table " TO STDOUT" + (when csv? " WITH CSV") + (when headers? " HEADER"))] + (with-open [out (ByteArrayOutputStream.)] + (.copyOut manager copy-query out) + (String. (.toByteArray out)))))) + +(defn copy-into! + "table-sql is the table name and columns for the COPY statement, + e.g. myschema.mytable(col1, col2). It should match the order of the + tuples exactly." + ([^java.sql.Connection conn + table-sql + values] + (let [manager (copy-manager conn) + vals (values->copy-rows-binary values)] + (with-open [stream (ByteArrayInputStream. vals)] + (.copyIn manager (str "COPY " (name table-sql) + " FROM STDIN WITH BINARY") stream)))) + ([^java.sql.Connection conn table cols values] + (let [table-spec (str (name table) + "(" + (str/join "," (map name cols)) + ")")] + (copy-into! conn table-spec values)))) diff --git a/src/clj_pgcopy/time.clj b/src/clj_pgcopy/time.clj new file mode 100644 index 0000000..1ce92d0 --- /dev/null +++ b/src/clj_pgcopy/time.clj @@ -0,0 +1,74 @@ +(ns clj-pgcopy.time + (:import (java.time LocalDateTime + LocalDate + ZoneOffset + ZonedDateTime + OffsetDateTime + Instant) + (java.util.concurrent TimeUnit))) + +(defn ^Long date-time->epoch-milli [^LocalDateTime dt] + (.. dt + (atOffset ZoneOffset/UTC) + toInstant + toEpochMilli)) + +(def epoch-distance + (- + (date-time->epoch-milli + ;; postgres epoch + (LocalDateTime/of 2000 1 1 0 0 0)) + (.toEpochMilli Instant/EPOCH))) + +;; 946684800000 + +;; The conversion is valid for any year 1583 CE onwards. +(defn java-epoch->postgres-epoch ^long [millis] + ;; returns microseconds + (long (* 1000 (- millis epoch-distance)))) + +(defn java-epoch->postgres-days [millis] + (->> (- millis epoch-distance) + (.toDays TimeUnit/MILLISECONDS) + int)) + +(defprotocol ToInstant + (-to-instant [self])) + +(extend-protocol ToInstant + LocalDateTime + (-to-instant [self] + (.. self + (atOffset ZoneOffset/UTC) + toInstant)) + LocalDate + (-to-instant [self] + (.. self + atStartOfDay + (atOffset ZoneOffset/UTC) + toInstant)) + Instant + (-to-instant [self] + self) + java.util.Date + (-to-instant [self] + (.toInstant self)) + java.sql.Date + (-to-instant [self] + (-to-instant + (.toLocalDate self))) + java.sql.Timestamp + (-to-instant [self] + (.toInstant self)) + ZonedDateTime + (-to-instant [self] + (.toInstant self)) + OffsetDateTime + (-to-instant [self] + (.toInstant self)) + nil + (-to-instant [_] + nil)) + +(defn to-instant [date-ish] + (-to-instant date-ish)) diff --git a/test/clj_pgcopy/core_test.clj b/test/clj_pgcopy/core_test.clj new file mode 100644 index 0000000..f4f22b9 --- /dev/null +++ b/test/clj_pgcopy/core_test.clj @@ -0,0 +1,236 @@ +(ns clj-pgcopy.core-test + (:require [clj-pgcopy.core :as copy] + [clj-pgcopy.time :as ptime] + [clojure.test :refer :all] + [clojure.java.jdbc :as jdbc]) + (:import (java.time LocalDateTime + LocalDate + Instant + ZoneId + ZoneOffset) + (org.postgresql.util PGInterval) + (org.postgresql.geometric PGbox + PGcircle + PGline + PGpath + PGpolygon + PGpoint) + (java.util TimeZone))) + +(def conn-spec "jdbc:postgresql://localhost:5432/test_pgcopy") + +(defprotocol IStringValue + (string-value [_])) + +(extend-protocol IStringValue + String + (string-value [this] this) + nil + (string-value [_] nil) + org.postgresql.jdbc.PgSQLXML + (string-value [this] (.getString this)) + Object + (string-value [this] (str this))) + +(use-fixtures :each + (fn use-utc [f] + (let [original (TimeZone/getDefault)] + (TimeZone/setDefault (TimeZone/getTimeZone "UTC")) + (f) + (TimeZone/setDefault original))) + (fn create-test-tables [f] + (jdbc/with-db-connection [conn conn-spec] + (jdbc/execute! conn "drop schema if exists copytest cascade") + (jdbc/execute! conn "create schema copytest") + (jdbc/execute! conn "create extension if not exists citext") + (jdbc/execute! conn (str "create table copytest.test(" + "c_int integer," + "c_bigint bigint," + "c_smallint smallint," + "c_text text," + "c_varchar varchar(32)," + "c_char char(2)," + "c_citext citext," + "c_text_array text[]," + "c_int_array int[]," + "c_double_array float8[]," + "c_uuid_array uuid[]," + "c_date date," + "c_ts timestamp," + "c_tstz timestamptz," + "c_boolean boolean," + "c_numeric numeric," + "c_decimal decimal(8,2)," + "c_float4 float4," + "c_float8 float8," + "c_uuid uuid," + "c_json json," + "c_jsonb jsonb," + "c_xml xml," + "c_interval interval," + "c_point point," + "c_line line," + "c_path path," + "c_box box," + "c_circle circle," + "c_polygon polygon," + "c_bytea bytea)")) + (f) + (jdbc/execute! conn "drop schema if exists copytest cascade")))) + +(deftest copy-columns-test + (let [row1 {:c_smallint (short 42) + :c_int (int 42) + :c_bigint (* 4 Integer/MAX_VALUE) + :c_text "Text" + :c_varchar "Varchar" + :c_char "ab" + :c_citext "CIText" + :c_text_array ["examples" "of" "text"] + :c_int_array [1 2 3 4 5] + :c_uuid_array [#uuid "cc7caf20-14e1-42de-a1e7-455d4f267111" + #uuid "cc7caf20-14e1-42de-a1e7-455d4f267112" + #uuid "cc7caf20-14e1-42de-a1e7-455d4f267113"] + :c_double_array [1.2 3.4 5.6] + :c_date (LocalDate/of 2018 1 2) + :c_ts (Instant/parse "2018-01-02T03:04:05.678Z") + :c_tstz (Instant/parse "2018-01-02T03:04:05.678Z") + :c_boolean true + :c_numeric 123456789.012345M + :c_decimal 4777.39M + :c_float4 (float 3.14159) + :c_float8 Math/PI + :c_uuid #uuid "cc7caf20-14e1-42de-a1e7-455d4f267111" + :c_json "{\"hello\": [1, 2, 3]}" + :c_jsonb (copy/->JsonB "{\"goodbye\": [3, 2, 1]}") + :c_xml "" + :c_interval (PGInterval. 1 2 3 4 5 6.789) + :c_point (PGpoint. "(1.2, -3.4)") + :c_line (PGline. "{1.2, -3.4, 5.6}") + :c_path (PGpath. "((1.2,-3.4),(-5.6,7.8),(9.0, 0.1))") + :c_polygon (PGpolygon. "((1.2,-3.4),(-5.6,7.8),(9.0, 0.1))") + :c_box (PGbox. "(-5.6,7.8),(1.2,-3.4)") + :c_circle (PGcircle. "<(1.2,-3.4), 9.1>") + :c_bytea "the byte array"} + columns (-> row1 keys vec) + data (vector row1 (into {} (map vector columns (repeat nil)))) + values (into [] (comp (map (fn [row] (-> row + (update :c_text_array #(into-array String %)) + (update :c_int_array int-array) + (update :c_double_array double-array) + (update :c_uuid_array #(into-array java.util.UUID %)) + (update :c_bytea #(if (string? %) + (.getBytes % "UTF-8") + %))))) + (map (apply juxt columns))) + data) + bytes->string #(when (and % (pos? (count %))) (String. %))] + (jdbc/with-db-connection [conn conn-spec] + (is (= (count data) + (copy/copy-into! (:connection conn) :copytest.test columns values))) + (let [results (->> (jdbc/query conn ["select * from copytest.test"]) + (map #(select-keys % columns )))] + (is (= (into [] + (comp + (map + (fn [row] + (-> row + (update :c_date ptime/to-instant) + (update :c_jsonb string-value)))) + ;;(map (apply juxt columns)) + ) + data) + (into [] + (comp + (map (fn [row] + (-> row + (update :c_date ptime/to-instant) + (update :c_ts ptime/to-instant) + (update :c_tstz ptime/to-instant) + (update :c_bytea bytes->string) + (update :c_text_array #(when % (not-empty (vec (.getArray %))))) + (update :c_int_array #(when % (not-empty (vec (.getArray %))))) + (update :c_double_array #(when % (not-empty (vec (.getArray %))))) + (update :c_uuid_array #(when % (not-empty (vec (.getArray %))))) + (update :c_citext string-value) + (update :c_json string-value) + (update :c_jsonb string-value) + (update :c_xml string-value)))) + ;;(map (apply juxt columns)) + ) + results))))))) + +(deftest datetime-inputs-test + (let [instant (Instant/parse "2018-01-02T03:04:05.678Z") + data [{;; ZonedDateTime -> timestamp + :c_ts (.atZone (LocalDateTime/ofInstant instant (ZoneId/of "Z")) + (ZoneId/of "Z")) + ;; ZonedDateTime -> timestamptz + :c_tstz (.atZone (LocalDateTime/ofInstant instant (ZoneId/of "Z")) + (ZoneId/of "Z"))} + + {;; OffsetDateTime -> timestamp + :c_ts (.atOffset (LocalDateTime/ofInstant instant (ZoneId/of "Z")) + ZoneOffset/UTC) + ;; OffsetDateTime -> timestamptz + :c_tstz (.atOffset (LocalDateTime/ofInstant instant (ZoneId/of "Z")) + ZoneOffset/UTC)} + + {;; Instant -> timestamp + :c_ts instant + ;; Instant -> timestamptz + :c_tstz instant} + + {;; java.util.Date -> timestamp + :c_ts (java.util.Date/from instant) + ;; java.util.Date -> timestamptz + :c_tstz (java.util.Date/from instant)} + + {;; java.sql.Timestamp -> timestamp + :c_ts (java.sql.Timestamp/from instant) + ;; java.sql.Timestamp -> timestamptz + :c_tstz (java.sql.Timestamp/from instant)}] + columns (-> data first keys vec) + values (into [] (map (apply juxt columns)) data) + fixup-row (fn [row] + (-> row + (update :c_tstz ptime/to-instant) + (update :c_ts ptime/to-instant))) + rows-xform (comp + (map fixup-row) + (map (apply juxt columns)))] + (jdbc/with-db-connection [conn conn-spec] + (is (= (count data) + (copy/copy-into! (:connection conn) :copytest.test columns values))) + (let [results (->> (jdbc/query conn ["select * from copytest.test"]) + (map #(select-keys % columns)))] + (is (= (into [] rows-xform data) + (into [] rows-xform results))))))) + +(deftest date-inputs-test + (let [data [;; LocalDate -> date + {:c_date (LocalDate/of 2018 1 2)} + ;; java.sql.Date -> date + {:c_date (java.sql.Date/valueOf (LocalDate/of 2018 1 2))} + ;; Far past dates + {:c_date (LocalDate/of 1980 1 2)} + {:c_date (java.sql.Date/valueOf (LocalDate/of 1582 10 4))} + {:c_date (java.sql.Date/valueOf (LocalDate/of 720 5 19))} + {:c_date (LocalDate/of 2 2 2)} + ;; Far future dates + {:c_date (LocalDate/of 4242 1 2)}] + columns (-> data first keys vec) + values (into [] (map (apply juxt columns)) data) + fixup-row (fn [row] + (-> row + (update :c_date ptime/to-instant))) + rows-xform (comp + (map fixup-row) + (map (apply juxt columns)))] + (jdbc/with-db-connection [conn conn-spec] + (is (= (count data) + (copy/copy-into! (:connection conn) :copytest.test columns values))) + (let [results (->> (jdbc/query conn ["select * from copytest.test"]) + (map #(select-keys % columns)))] + (is (= (into [] rows-xform data) + (into [] rows-xform results)))))))