From 70f23f4711fbc6773f52d123752e9cbf6ea3f73f Mon Sep 17 00:00:00 2001 From: Lukasz Samson Date: Fri, 25 Sep 2026 22:46:12 +0200 Subject: [PATCH 1/7] Map insert_all select fields through destination schema --- lib/ecto/repo/schema.ex | 42 +++++++++++++++++--- test/ecto/repo_test.exs | 86 +++++++++++++++++++++++++++++++++++++++++ 2 files changed, 122 insertions(+), 6 deletions(-) diff --git a/lib/ecto/repo/schema.ex b/lib/ecto/repo/schema.ex index 4d32d2f95b..54c81377dc 100644 --- a/lib/ecto/repo/schema.ex +++ b/lib/ecto/repo/schema.ex @@ -179,13 +179,13 @@ defmodule Ecto.Repo.Schema do {updated_fields, updated_set} = Enum.map_reduce(args, MapSet.new(), fn {field, _}, set -> dumped_field = insert_all_select_dump!(field, dumper) - {dumped_field, MapSet.put(set, dumped_field)} + {dumped_field, MapSet.put(set, field)} end) + source_fields = Enum.take(fields, length(fields) - length(args)) + unchanged_fields = - for {{:., _, [{:&, _, [^ix]}, field]}, [], []} = expr <- fields, - not MapSet.member?(updated_set, field), - do: insert_all_select_dump!(expr) + insert_all_source_fields(query, ix, source_fields, updated_set, dumper) unchanged_fields ++ updated_fields @@ -195,8 +195,8 @@ defmodule Ecto.Repo.Schema do %Ecto.Query.SelectExpr{take: %{^ix => {_fun, fields}}} -> Enum.map(fields, &insert_all_select_dump!(&1, dumper)) - %Ecto.Query.SelectExpr{expr: {:&, _, [_ix]}, fields: fields} -> - Enum.map(fields, &insert_all_select_dump!(&1)) + %Ecto.Query.SelectExpr{expr: {:&, _, [ix]}, fields: fields} -> + insert_all_source_fields(query, ix, fields, MapSet.new(), dumper) _ -> raise ArgumentError, """ @@ -344,6 +344,36 @@ defmodule Ecto.Repo.Schema do end end + defp insert_all_source_fields(_query, _ix, fields, _updated_set, nil) do + Enum.map(fields, &insert_all_select_dump!/1) + end + + defp insert_all_source_fields(query, ix, fields, updated_set, dumper) do + source_fields = + case elem(query.sources, ix) do + {_, schema, _} when is_atom(schema) and not is_nil(schema) -> + case query.select.take do + %{^ix => {_fun, selected_fields}} -> selected_fields + _ -> schema.__schema__(:query_fields) + end + |> Enum.reject(&MapSet.member?(updated_set, &1)) + + _ -> + Enum.map(fields, &insert_all_select_dump!/1) + end + + if length(source_fields) != length(fields) do + raise ArgumentError, + "cannot generate a fields list for insert_all from the given source query: " <> + inspect(query) + end + + Enum.zip_with(source_fields, fields, fn field, expr -> + insert_all_select_dump!(expr) + insert_all_select_dump!(field, dumper) + end) + end + defp insert_all_select_dump!(field, dumper) when is_atom(field) do case dumper do %{^field => {source, _, writable}} when writable != :never -> diff --git a/test/ecto/repo_test.exs b/test/ecto/repo_test.exs index bd6248f711..cb80cb089a 100644 --- a/test/ecto/repo_test.exs +++ b/test/ecto/repo_test.exs @@ -211,6 +211,46 @@ defmodule Ecto.RepoTest do end end + defmodule InsertSelectSource do + use Ecto.Schema + + @primary_key false + schema "insert_select_source" do + field :name, :string + field :value, :string + end + end + + defmodule InsertSelectMappedSource do + use Ecto.Schema + + @primary_key false + schema "insert_select_source" do + field :name, :string, source: :source_name + field :value, :string + end + end + + defmodule InsertSelectRenamed do + use Ecto.Schema + + @primary_key false + schema "insert_select_renamed" do + field :name, :string, source: :renamed_name + field :value, :string + end + end + + defmodule InsertSelectReadOnly do + use Ecto.Schema + + @primary_key false + schema "insert_select_read_only" do + field :name, :string, writable: :never + field :value, :string + end + end + test "defines child_spec/1" do assert TestRepo.child_spec([]) == %{ id: TestRepo, @@ -765,6 +805,52 @@ defmodule Ecto.RepoTest do assert header == [:id, :x, :yyy, :z, :array, :map] end + test "maps full source fields through the destination schema" do + query = from s in InsertSelectSource, select: s + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + + test "keeps source columns when the destination has no schema" do + query = from s in InsertSelectMappedSource, select: s + TestRepo.insert_all("insert_select_renamed", query) + + assert_received {:insert_all, %{header: [:source_name, :value]}, + {%Ecto.Query{}, _params}} + end + + test "rejects full source fields that are unwritable in the destination" do + query = from s in InsertSelectSource, select: s + + assert_raise ArgumentError, + "cannot select unwritable field `:name` for insert_all", + fn -> TestRepo.insert_all(InsertSelectReadOnly, query) end + end + + test "maps unchanged map update fields through the destination schema" do + query = from s in InsertSelectMappedSource, select: %{s | value: s.value} + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + + test "does not include an overwritten source field twice when columns differ" do + query = from s in InsertSelectMappedSource, select: %{s | name: s.name} + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:value, :renamed_name]}, + {%Ecto.Query{}, _params}} + end + + test "rejects unchanged map update fields that are unwritable in the destination" do + query = from s in InsertSelectMappedSource, select: %{s | value: s.value} + + assert_raise ArgumentError, + "cannot select unwritable field `:name` for insert_all", + fn -> TestRepo.insert_all(InsertSelectReadOnly, query) end + end + test "takes query selecting on source with join" do query = from p in MyParent, join: a in MySchemaWithAssoc, on: true, select: a TestRepo.insert_all(MySchemaWithAssoc, query) From 3e6d2caed91b97911ebaada7bb8518a8f258dda7 Mon Sep 17 00:00:00 2001 From: Lukasz Samson Date: Fri, 25 Sep 2026 23:21:28 +0200 Subject: [PATCH 2/7] Validate insert_all source projections and writable destinations --- lib/ecto/repo/schema.ex | 66 ++++++++++++++++++++--------------------- test/ecto/repo_test.exs | 32 ++++++++++++++++++++ 2 files changed, 65 insertions(+), 33 deletions(-) diff --git a/lib/ecto/repo/schema.ex b/lib/ecto/repo/schema.ex index 54c81377dc..e4ba510ca0 100644 --- a/lib/ecto/repo/schema.ex +++ b/lib/ecto/repo/schema.ex @@ -199,24 +199,7 @@ defmodule Ecto.Repo.Schema do insert_all_source_fields(query, ix, fields, MapSet.new(), dumper) _ -> - raise ArgumentError, """ - cannot generate a fields list for insert_all from the given source query: - - #{inspect(query)} - - The select clause must be one of the following: - - * A single `map/2` or several `map/2` expressions combined with `select_merge` - * A single `struct/2` or several `struct/2` expressions combined with `select_merge` - * A source such as `p` in the query `from p in Post` - * A single literal map or several literal maps combined with `select_merge`. If - combining several literal maps, there cannot be any query interpolations - except in the last `select_merge`. Consider using `Ecto.Query.exclude/2` - to rebuild the select expression from scratch if you need multiple `select_merge` - statements with interpolations - - All keys must exist in the schema that is being inserted into - """ + insert_all_select_error!(query) end counter = fn -> length(dump_params) end @@ -336,16 +319,10 @@ defmodule Ecto.Repo.Schema do {rows, Enum.reverse(cast_params), counter} end - defp insert_all_select_dump!({{:., dot_meta, [{:&, _, [_]}, field]}, [], []}) do - if dot_meta[:writable] == :never do - raise ArgumentError, "cannot select unwritable field `#{inspect(field)}` for insert_all" - else - field - end - end + defp insert_all_select_source!({{:., _, [{:&, _, [_]}, field]}, [], []}), do: field defp insert_all_source_fields(_query, _ix, fields, _updated_set, nil) do - Enum.map(fields, &insert_all_select_dump!/1) + Enum.map(fields, &insert_all_select_source!/1) end defp insert_all_source_fields(query, ix, fields, updated_set, dumper) do @@ -356,24 +333,47 @@ defmodule Ecto.Repo.Schema do %{^ix => {_fun, selected_fields}} -> selected_fields _ -> schema.__schema__(:query_fields) end + |> Enum.filter(&is_atom/1) |> Enum.reject(&MapSet.member?(updated_set, &1)) _ -> - Enum.map(fields, &insert_all_select_dump!/1) + Enum.map(fields, &insert_all_select_source!/1) end if length(source_fields) != length(fields) do - raise ArgumentError, - "cannot generate a fields list for insert_all from the given source query: " <> - inspect(query) + insert_all_select_error!(query) end - Enum.zip_with(source_fields, fields, fn field, expr -> - insert_all_select_dump!(expr) - insert_all_select_dump!(field, dumper) + Enum.zip_with(source_fields, fields, fn + field, {{:., _, [{:&, _, [^ix]}, _]}, [], []} -> + insert_all_select_dump!(field, dumper) + + _, _ -> + insert_all_select_error!(query) end) end + defp insert_all_select_error!(query) do + raise ArgumentError, """ + cannot generate a fields list for insert_all from the given source query: + + #{inspect(query)} + + The select clause must be one of the following: + + * A single `map/2` or several `map/2` expressions combined with `select_merge` + * A single `struct/2` or several `struct/2` expressions combined with `select_merge` + * A source such as `p` in the query `from p in Post` + * A single literal map or several literal maps combined with `select_merge`. If + combining several literal maps, there cannot be any query interpolations + except in the last `select_merge`. Consider using `Ecto.Query.exclude/2` + to rebuild the select expression from scratch if you need multiple `select_merge` + statements with interpolations + + All keys must exist in the schema that is being inserted into + """ + end + defp insert_all_select_dump!(field, dumper) when is_atom(field) do case dumper do %{^field => {source, _, writable}} when writable != :never -> diff --git a/test/ecto/repo_test.exs b/test/ecto/repo_test.exs index cb80cb089a..2da869d183 100644 --- a/test/ecto/repo_test.exs +++ b/test/ecto/repo_test.exs @@ -231,6 +231,16 @@ defmodule Ecto.RepoTest do end end + defmodule InsertSelectReadOnlySource do + use Ecto.Schema + + @primary_key false + schema "insert_select_source" do + field :name, :string, writable: :never + field :value, :string + end + end + defmodule InsertSelectRenamed do use Ecto.Schema @@ -828,6 +838,13 @@ defmodule Ecto.RepoTest do fn -> TestRepo.insert_all(InsertSelectReadOnly, query) end end + test "can read an unwritable source field into a writable destination" do + query = from s in InsertSelectReadOnlySource, select: s + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + test "maps unchanged map update fields through the destination schema" do query = from s in InsertSelectMappedSource, select: %{s | value: s.value} TestRepo.insert_all(InsertSelectRenamed, query) @@ -835,6 +852,13 @@ defmodule Ecto.RepoTest do assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} end + test "maps unchanged fields from a map subset through the destination schema" do + query = from s in InsertSelectMappedSource, select: %{map(s, [:name]) | value: s.value} + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + test "does not include an overwritten source field twice when columns differ" do query = from s in InsertSelectMappedSource, select: %{s | name: s.name} TestRepo.insert_all(InsertSelectRenamed, query) @@ -851,6 +875,14 @@ defmodule Ecto.RepoTest do fn -> TestRepo.insert_all(InsertSelectReadOnly, query) end end + test "rejects map updates whose values expand to multiple select fields" do + query = from s in InsertSelectSource, select: %{s | value: {s.name, s.value}} + + assert_raise ArgumentError, + ~r/cannot generate a fields list for insert_all from the given source query:/, + fn -> TestRepo.insert_all(InsertSelectRenamed, query) end + end + test "takes query selecting on source with join" do query = from p in MyParent, join: a in MySchemaWithAssoc, on: true, select: a TestRepo.insert_all(MySchemaWithAssoc, query) From 3b8d07c1f3a7f79741cfe57aa1dd6f4a3abac235 Mon Sep 17 00:00:00 2001 From: Lukasz Samson Date: Sun, 27 Sep 2026 07:31:42 +0200 Subject: [PATCH 3/7] Verify insert_all source projection against planner fields --- lib/ecto/repo/schema.ex | 93 +++++++++++++++++++++++++++++------------ test/ecto/repo_test.exs | 74 ++++++++++++++++++++++++++++++++ 2 files changed, 140 insertions(+), 27 deletions(-) diff --git a/lib/ecto/repo/schema.ex b/lib/ecto/repo/schema.ex index e4ba510ca0..b45b125d0d 100644 --- a/lib/ecto/repo/schema.ex +++ b/lib/ecto/repo/schema.ex @@ -182,10 +182,8 @@ defmodule Ecto.Repo.Schema do {dumped_field, MapSet.put(set, field)} end) - source_fields = Enum.take(fields, length(fields) - length(args)) - unchanged_fields = - insert_all_source_fields(query, ix, source_fields, updated_set, dumper) + insert_all_source_fields(query, ix, fields, updated_set, length(args), dumper) unchanged_fields ++ updated_fields @@ -196,7 +194,7 @@ defmodule Ecto.Repo.Schema do Enum.map(fields, &insert_all_select_dump!(&1, dumper)) %Ecto.Query.SelectExpr{expr: {:&, _, [ix]}, fields: fields} -> - insert_all_source_fields(query, ix, fields, MapSet.new(), dumper) + insert_all_source_fields(query, ix, fields, MapSet.new(), 0, dumper) _ -> insert_all_select_error!(query) @@ -319,40 +317,74 @@ defmodule Ecto.Repo.Schema do {rows, Enum.reverse(cast_params), counter} end - defp insert_all_select_source!({{:., _, [{:&, _, [_]}, field]}, [], []}), do: field - - defp insert_all_source_fields(_query, _ix, fields, _updated_set, nil) do - Enum.map(fields, &insert_all_select_source!/1) - end - - defp insert_all_source_fields(query, ix, fields, updated_set, dumper) do + defp insert_all_source_fields(query, ix, fields, updated_set, updated_count, dumper) do + {source_fields, source_dumper} = insert_all_source_projection(query, ix) source_fields = - case elem(query.sources, ix) do - {_, schema, _} when is_atom(schema) and not is_nil(schema) -> - case query.select.take do - %{^ix => {_fun, selected_fields}} -> selected_fields - _ -> schema.__schema__(:query_fields) - end - |> Enum.filter(&is_atom/1) - |> Enum.reject(&MapSet.member?(updated_set, &1)) + source_fields + |> Enum.filter(&is_atom/1) + |> Enum.reject(&MapSet.member?(updated_set, &1)) - _ -> - Enum.map(fields, &insert_all_select_source!/1) - end + {source_exprs, updated_exprs} = Enum.split(fields, length(source_fields)) - if length(source_fields) != length(fields) do + if length(source_exprs) != length(source_fields) or length(updated_exprs) != updated_count do insert_all_select_error!(query) end - Enum.zip_with(source_fields, fields, fn - field, {{:., _, [{:&, _, [^ix]}, _]}, [], []} -> - insert_all_select_dump!(field, dumper) + # The planner expands source fields to physical columns in this order. + Enum.zip_with(source_fields, source_exprs, fn + field, {{:., _, [{:&, _, [^ix]}, source]}, [], []} -> + {expected_source, _, _} = Map.get(source_dumper, field, {field, :any, :always}) + + if source != expected_source do + insert_all_select_error!(query) + end + + if dumper, do: insert_all_select_dump!(field, dumper), else: source _, _ -> insert_all_select_error!(query) end) end + defp insert_all_source_projection(query, ix) do + source = elem(query.sources, ix) + + fields = + case query.select.take do + %{^ix => {_fun, fields}} -> fields + _ -> nil + end + + projection = + case source do + {_, schema, _} when is_atom(schema) and not is_nil(schema) -> + {fields || schema.__schema__(:query_fields), schema.__schema__(:dump)} + + %Ecto.SubQuery{select: {:source, _, _, types}} -> + {fields || Keyword.keys(types), %{}} + + %Ecto.SubQuery{select: {:struct, _, types}} -> + {fields || Keyword.keys(types), %{}} + + %Ecto.SubQuery{select: {:map, types}} -> + {fields || Keyword.keys(types), %{}} + + {{:fragment, meta, _}, nil, _} -> + {fields || meta[:column_names], %{}} + + {:values, _, [types, _]} -> + {fields || Keyword.keys(types), %{}} + + _ -> + {fields, %{}} + end + + case projection do + {nil, _} -> insert_all_select_error!(query) + projection -> projection + end + end + defp insert_all_select_error!(query) do raise ArgumentError, """ cannot generate a fields list for insert_all from the given source query: @@ -379,14 +411,21 @@ defmodule Ecto.Repo.Schema do %{^field => {source, _, writable}} when writable != :never -> source - %{} -> + %{^field => {_, _, :never}} -> raise ArgumentError, "cannot select unwritable field `#{inspect(field)}` for insert_all" + %{} -> + raise ArgumentError, "cannot select unknown field `#{inspect(field)}` for insert_all" + nil -> field end end + defp insert_all_select_dump!(field, _dumper) do + raise ArgumentError, "cannot select non-atom field `#{inspect(field)}` for insert_all" + end + defp autogenerate_id(nil, fields, header, _adapter) do {fields, header} end diff --git a/test/ecto/repo_test.exs b/test/ecto/repo_test.exs index 2da869d183..5a325268a1 100644 --- a/test/ecto/repo_test.exs +++ b/test/ecto/repo_test.exs @@ -261,6 +261,30 @@ defmodule Ecto.RepoTest do end end + defmodule InsertSelectDisjointSource do + use Ecto.Schema + + @primary_key false + schema "insert_select_disjoint_source" do + field :b, :string, source: :src_b + field :c, :string, source: :src_c + field :a, :string, source: :src_a + field :only_source, :string + end + end + + defmodule InsertSelectDisjointDestination do + use Ecto.Schema + + @primary_key false + schema "insert_select_disjoint_destination" do + field :only_destination, :string + field :a, :string + field :c, :string, source: :dst_c + field :b, :string, source: :dst_b + end + end + test "defines child_spec/1" do assert TestRepo.child_spec([]) == %{ id: TestRepo, @@ -883,6 +907,56 @@ defmodule Ecto.RepoTest do fn -> TestRepo.insert_all(InsertSelectRenamed, query) end end + test "maps a reordered source subset independently of destination field order" do + query = + from s in InsertSelectDisjointSource, + select: %{map(s, [:b, :c, :a]) | a: fragment("'x'")} + + TestRepo.insert_all(InsertSelectDisjointDestination, query) + + assert_received {:insert_all, %{header: [:dst_b, :dst_c, :a]}, + {%Ecto.Query{select: %{fields: fields}}, _params}} + + assert [ + {{:., _, [{:&, _, [0]}, :src_b]}, [], []}, + {{:., _, [{:&, _, [0]}, :src_c]}, [], []}, + {:fragment, _, _} + ] = fields + end + + test "allows destination-only map updates and reports source-only fields" do + query = + from s in InsertSelectDisjointSource, + select: %{map(s, [:b, :c]) | a: s.c, only_destination: s.b} + + TestRepo.insert_all(InsertSelectDisjointDestination, query) + + assert_received {:insert_all, %{header: [:dst_b, :dst_c, :a, :only_destination]}, + {%Ecto.Query{}, _params}} + + query = from s in InsertSelectDisjointSource, select: s + + assert_raise ArgumentError, + "cannot select unknown field `:only_source` for insert_all", + fn -> TestRepo.insert_all(InsertSelectDisjointDestination, query) end + end + + test "rejects whole schemaless bindings with a select error" do + query = from s in "insert_select_source", select: s + + assert_raise ArgumentError, + ~r/cannot generate a fields list for insert_all from the given source query:/, + fn -> TestRepo.insert_all(InsertSelectRenamed, query) end + end + + test "rejects non-atom insert select keys with a clear error" do + query = from s in InsertSelectSource, select: %{"name" => s.name} + + assert_raise ArgumentError, + "cannot select non-atom field `\"name\"` for insert_all", + fn -> TestRepo.insert_all(InsertSelectRenamed, query) end + end + test "takes query selecting on source with join" do query = from p in MyParent, join: a in MySchemaWithAssoc, on: true, select: a TestRepo.insert_all(MySchemaWithAssoc, query) From 592e193250141d3746cd0021dec56c8ac7828768 Mon Sep 17 00:00:00 2001 From: Lukasz Samson Date: Sun, 27 Sep 2026 08:18:41 +0200 Subject: [PATCH 4/7] Use planner fields for non-schema insert selections --- lib/ecto/repo/schema.ex | 90 ++++++++++++++++++----------------------- test/ecto/repo_test.exs | 47 ++++++++++++++++++++- 2 files changed, 86 insertions(+), 51 deletions(-) diff --git a/lib/ecto/repo/schema.ex b/lib/ecto/repo/schema.ex index b45b125d0d..fcfc2b035f 100644 --- a/lib/ecto/repo/schema.ex +++ b/lib/ecto/repo/schema.ex @@ -318,73 +318,63 @@ defmodule Ecto.Repo.Schema do end defp insert_all_source_fields(query, ix, fields, updated_set, updated_count, dumper) do - {source_fields, source_dumper} = insert_all_source_projection(query, ix) - source_fields = - source_fields - |> Enum.filter(&is_atom/1) - |> Enum.reject(&MapSet.member?(updated_set, &1)) - - {source_exprs, updated_exprs} = Enum.split(fields, length(source_fields)) - - if length(source_exprs) != length(source_fields) or length(updated_exprs) != updated_count do - insert_all_select_error!(query) - end + case elem(query.sources, ix) do + {_, schema, _} when is_atom(schema) and not is_nil(schema) -> + source_fields = + case query.select.take do + %{^ix => {_fun, selected_fields}} -> selected_fields + _ -> schema.__schema__(:query_fields) + end + |> Enum.filter(&is_atom/1) + |> Enum.reject(&MapSet.member?(updated_set, &1)) - # The planner expands source fields to physical columns in this order. - Enum.zip_with(source_fields, source_exprs, fn - field, {{:., _, [{:&, _, [^ix]}, source]}, [], []} -> - {expected_source, _, _} = Map.get(source_dumper, field, {field, :any, :always}) + {source_exprs, updated_exprs} = Enum.split(fields, length(source_fields)) - if source != expected_source do + if length(source_exprs) != length(source_fields) or length(updated_exprs) != updated_count do insert_all_select_error!(query) end - if dumper, do: insert_all_select_dump!(field, dumper), else: source + source_dumper = schema.__schema__(:dump) - _, _ -> - insert_all_select_error!(query) - end) - end + Enum.zip_with(source_fields, source_exprs, fn field, expr -> + source = insert_all_source_field!(query, ix, expr) + {expected_source, _, _} = Map.get(source_dumper, field, {field, :any, :always}) - defp insert_all_source_projection(query, ix) do - source = elem(query.sources, ix) - - fields = - case query.select.take do - %{^ix => {_fun, fields}} -> fields - _ -> nil - end + if source != expected_source do + insert_all_select_error!(query) + end - projection = - case source do - {_, schema, _} when is_atom(schema) and not is_nil(schema) -> - {fields || schema.__schema__(:query_fields), schema.__schema__(:dump)} + if dumper, do: insert_all_select_dump!(field, dumper), else: source + end) - %Ecto.SubQuery{select: {:source, _, _, types}} -> - {fields || Keyword.keys(types), %{}} + _ -> + {source_exprs, updated_exprs} = split_updated_fields(fields, updated_count) - %Ecto.SubQuery{select: {:struct, _, types}} -> - {fields || Keyword.keys(types), %{}} + if length(updated_exprs) != updated_count do + insert_all_select_error!(query) + end - %Ecto.SubQuery{select: {:map, types}} -> - {fields || Keyword.keys(types), %{}} + Enum.map(source_exprs, fn expr -> + field = insert_all_source_field!(query, ix, expr) - {{:fragment, meta, _}, nil, _} -> - {fields || meta[:column_names], %{}} + if MapSet.member?(updated_set, field) do + insert_all_select_error!(query) + end - {:values, _, [types, _]} -> - {fields || Keyword.keys(types), %{}} + insert_all_select_dump!(field, dumper) + end) + end + end - _ -> - {fields, %{}} - end + defp split_updated_fields(fields, 0), do: {fields, []} + defp split_updated_fields(fields, count), do: Enum.split(fields, -count) - case projection do - {nil, _} -> insert_all_select_error!(query) - projection -> projection - end + defp insert_all_source_field!(query, ix, {{:., _, [{:&, _, [expr_ix]}, field]}, [], []}) do + if expr_ix == ix, do: field, else: insert_all_select_error!(query) end + defp insert_all_source_field!(query, _ix, _expr), do: insert_all_select_error!(query) + defp insert_all_select_error!(query) do raise ArgumentError, """ cannot generate a fields list for insert_all from the given source query: diff --git a/test/ecto/repo_test.exs b/test/ecto/repo_test.exs index 5a325268a1..76ebee8788 100644 --- a/test/ecto/repo_test.exs +++ b/test/ecto/repo_test.exs @@ -846,6 +846,49 @@ defmodule Ecto.RepoTest do assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} end + test "maps fragment columns through the destination schema" do + query = + from f in fragment("select 1 as name, 2 as value", columns: [:name, :value]), select: f + + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + + test "maps subquery fields through the destination schema" do + inner = from s in InsertSelectSource, select: %{name: s.name, value: s.value} + query = from s in subquery(inner), select: s + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + + test "rejects a map update that repeats a subquery field" do + inner = from s in InsertSelectSource, select: %{name: s.name, value: s.value} + query = from s in subquery(inner), select: %{s | value: "new"} + + assert_raise ArgumentError, + ~r/cannot generate a fields list for insert_all from the given source query:/, + fn -> TestRepo.insert_all(InsertSelectRenamed, query) end + end + + test "maps values fields through the destination schema" do + query = + from v in values([%{name: "n", value: "v"}], %{name: :string, value: :string}), + select: %{v | value: "new"} + + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + + test "maps a joined binding through the destination schema" do + query = from x in "other", join: s in InsertSelectSource, on: true, select: s + TestRepo.insert_all(InsertSelectRenamed, query) + + assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + end + test "keeps source columns when the destination has no schema" do query = from s in InsertSelectMappedSource, select: s TestRepo.insert_all("insert_select_renamed", query) @@ -924,7 +967,7 @@ defmodule Ecto.RepoTest do ] = fields end - test "allows destination-only map updates and reports source-only fields" do + test "allows destination-only map updates" do query = from s in InsertSelectDisjointSource, select: %{map(s, [:b, :c]) | a: s.c, only_destination: s.b} @@ -933,7 +976,9 @@ defmodule Ecto.RepoTest do assert_received {:insert_all, %{header: [:dst_b, :dst_c, :a, :only_destination]}, {%Ecto.Query{}, _params}} + end + test "reports source-only fields" do query = from s in InsertSelectDisjointSource, select: s assert_raise ArgumentError, From edde18d01e985c862e5320ee01e83d89e559592f Mon Sep 17 00:00:00 2001 From: Lukasz Samson Date: Sun, 27 Sep 2026 09:04:08 +0200 Subject: [PATCH 5/7] Clarify insert_all projection alignment tests --- lib/ecto/repo/schema.ex | 2 ++ test/ecto/repo_test.exs | 18 ++++++++++++++++-- 2 files changed, 18 insertions(+), 2 deletions(-) diff --git a/lib/ecto/repo/schema.ex b/lib/ecto/repo/schema.ex index fcfc2b035f..086853a3af 100644 --- a/lib/ecto/repo/schema.ex +++ b/lib/ecto/repo/schema.ex @@ -318,6 +318,8 @@ defmodule Ecto.Repo.Schema do end defp insert_all_source_fields(query, ix, fields, updated_set, updated_count, dumper) do + # Keep insert headers aligned with the planner's SELECT fields by checking the + # physical source column before mapping each logical field to its destination. case elem(query.sources, ix) do {_, schema, _} when is_atom(schema) and not is_nil(schema) -> source_fields = diff --git a/test/ecto/repo_test.exs b/test/ecto/repo_test.exs index 76ebee8788..cbf01b6fef 100644 --- a/test/ecto/repo_test.exs +++ b/test/ecto/repo_test.exs @@ -843,7 +843,13 @@ defmodule Ecto.RepoTest do query = from s in InsertSelectSource, select: s TestRepo.insert_all(InsertSelectRenamed, query) - assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + assert_received {:insert_all, %{header: [:renamed_name, :value]}, + {%Ecto.Query{select: %{fields: fields}}, _params}} + + assert [ + {{:., _, [{:&, _, [0]}, :name]}, [], []}, + {{:., _, [{:&, _, [0]}, :value]}, [], []} + ] = fields end test "maps fragment columns through the destination schema" do @@ -867,6 +873,8 @@ defmodule Ecto.RepoTest do inner = from s in InsertSelectSource, select: %{name: s.name, value: s.value} query = from s in subquery(inner), select: %{s | value: "new"} + # This query shape still projects the overwritten field from the subquery. + # Reject the duplicate projection; this does not add support for subquery map updates. assert_raise ArgumentError, ~r/cannot generate a fields list for insert_all from the given source query:/, fn -> TestRepo.insert_all(InsertSelectRenamed, query) end @@ -916,7 +924,13 @@ defmodule Ecto.RepoTest do query = from s in InsertSelectMappedSource, select: %{s | value: s.value} TestRepo.insert_all(InsertSelectRenamed, query) - assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} + assert_received {:insert_all, %{header: [:renamed_name, :value]}, + {%Ecto.Query{select: %{fields: fields}}, _params}} + + assert [ + {{:., _, [{:&, _, [0]}, :source_name]}, [], []}, + {{:., _, [{:&, _, [0]}, :value]}, [], []} + ] = fields end test "maps unchanged fields from a map subset through the destination schema" do From 88b73efbddc19470ce92e88bab76ab2c31305e48 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Valim?= Date: Sun, 27 Sep 2026 10:37:05 +0200 Subject: [PATCH 6/7] Reuse fixtures when possible --- test/ecto/repo_test.exs | 78 +++++++++-------------------------------- 1 file changed, 16 insertions(+), 62 deletions(-) diff --git a/test/ecto/repo_test.exs b/test/ecto/repo_test.exs index cbf01b6fef..6ce893eff5 100644 --- a/test/ecto/repo_test.exs +++ b/test/ecto/repo_test.exs @@ -211,36 +211,6 @@ defmodule Ecto.RepoTest do end end - defmodule InsertSelectSource do - use Ecto.Schema - - @primary_key false - schema "insert_select_source" do - field :name, :string - field :value, :string - end - end - - defmodule InsertSelectMappedSource do - use Ecto.Schema - - @primary_key false - schema "insert_select_source" do - field :name, :string, source: :source_name - field :value, :string - end - end - - defmodule InsertSelectReadOnlySource do - use Ecto.Schema - - @primary_key false - schema "insert_select_source" do - field :name, :string, writable: :never - field :value, :string - end - end - defmodule InsertSelectRenamed do use Ecto.Schema @@ -256,7 +226,7 @@ defmodule Ecto.RepoTest do @primary_key false schema "insert_select_read_only" do - field :name, :string, writable: :never + field :name, :string, source: :source_name, writable: :never field :value, :string end end @@ -839,15 +809,15 @@ defmodule Ecto.RepoTest do assert header == [:id, :x, :yyy, :z, :array, :map] end - test "maps full source fields through the destination schema" do - query = from s in InsertSelectSource, select: s + test "maps read-only source fields through the writable destination schema" do + query = from s in InsertSelectReadOnly, select: s TestRepo.insert_all(InsertSelectRenamed, query) assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{select: %{fields: fields}}, _params}} assert [ - {{:., _, [{:&, _, [0]}, :name]}, [], []}, + {{:., _, [{:&, _, [0]}, :source_name]}, [], []}, {{:., _, [{:&, _, [0]}, :value]}, [], []} ] = fields end @@ -862,7 +832,7 @@ defmodule Ecto.RepoTest do end test "maps subquery fields through the destination schema" do - inner = from s in InsertSelectSource, select: %{name: s.name, value: s.value} + inner = from s in InsertSelectReadOnly, select: %{name: s.name, value: s.value} query = from s in subquery(inner), select: s TestRepo.insert_all(InsertSelectRenamed, query) @@ -870,7 +840,7 @@ defmodule Ecto.RepoTest do end test "rejects a map update that repeats a subquery field" do - inner = from s in InsertSelectSource, select: %{name: s.name, value: s.value} + inner = from s in InsertSelectReadOnly, select: %{name: s.name, value: s.value} query = from s in subquery(inner), select: %{s | value: "new"} # This query shape still projects the overwritten field from the subquery. @@ -891,37 +861,29 @@ defmodule Ecto.RepoTest do end test "maps a joined binding through the destination schema" do - query = from x in "other", join: s in InsertSelectSource, on: true, select: s + query = from x in "other", join: s in InsertSelectReadOnly, on: true, select: s TestRepo.insert_all(InsertSelectRenamed, query) assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} end test "keeps source columns when the destination has no schema" do - query = from s in InsertSelectMappedSource, select: s + query = from s in InsertSelectReadOnly, select: s TestRepo.insert_all("insert_select_renamed", query) - assert_received {:insert_all, %{header: [:source_name, :value]}, - {%Ecto.Query{}, _params}} + assert_received {:insert_all, %{header: [:source_name, :value]}, {%Ecto.Query{}, _params}} end test "rejects full source fields that are unwritable in the destination" do - query = from s in InsertSelectSource, select: s + query = from s in InsertSelectRenamed, select: s assert_raise ArgumentError, "cannot select unwritable field `:name` for insert_all", fn -> TestRepo.insert_all(InsertSelectReadOnly, query) end end - test "can read an unwritable source field into a writable destination" do - query = from s in InsertSelectReadOnlySource, select: s - TestRepo.insert_all(InsertSelectRenamed, query) - - assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} - end - test "maps unchanged map update fields through the destination schema" do - query = from s in InsertSelectMappedSource, select: %{s | value: s.value} + query = from s in InsertSelectReadOnly, select: %{s | value: s.value} TestRepo.insert_all(InsertSelectRenamed, query) assert_received {:insert_all, %{header: [:renamed_name, :value]}, @@ -933,23 +895,15 @@ defmodule Ecto.RepoTest do ] = fields end - test "maps unchanged fields from a map subset through the destination schema" do - query = from s in InsertSelectMappedSource, select: %{map(s, [:name]) | value: s.value} - TestRepo.insert_all(InsertSelectRenamed, query) - - assert_received {:insert_all, %{header: [:renamed_name, :value]}, {%Ecto.Query{}, _params}} - end - test "does not include an overwritten source field twice when columns differ" do - query = from s in InsertSelectMappedSource, select: %{s | name: s.name} + query = from s in InsertSelectReadOnly, select: %{s | name: s.name} TestRepo.insert_all(InsertSelectRenamed, query) - assert_received {:insert_all, %{header: [:value, :renamed_name]}, - {%Ecto.Query{}, _params}} + assert_received {:insert_all, %{header: [:value, :renamed_name]}, {%Ecto.Query{}, _params}} end test "rejects unchanged map update fields that are unwritable in the destination" do - query = from s in InsertSelectMappedSource, select: %{s | value: s.value} + query = from s in InsertSelectRenamed, select: %{s | value: s.value} assert_raise ArgumentError, "cannot select unwritable field `:name` for insert_all", @@ -957,7 +911,7 @@ defmodule Ecto.RepoTest do end test "rejects map updates whose values expand to multiple select fields" do - query = from s in InsertSelectSource, select: %{s | value: {s.name, s.value}} + query = from s in InsertSelectReadOnly, select: %{s | value: {s.name, s.value}} assert_raise ArgumentError, ~r/cannot generate a fields list for insert_all from the given source query:/, @@ -1009,7 +963,7 @@ defmodule Ecto.RepoTest do end test "rejects non-atom insert select keys with a clear error" do - query = from s in InsertSelectSource, select: %{"name" => s.name} + query = from s in InsertSelectReadOnly, select: %{"name" => s.name} assert_raise ArgumentError, "cannot select non-atom field `\"name\"` for insert_all", From d195bcc614526e33a602d648a126338fc84e1564 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Valim?= Date: Sun, 27 Sep 2026 11:01:36 +0200 Subject: [PATCH 7/7] Simplify and optimize field rendering logic --- lib/ecto/repo/schema.ex | 94 ++++++++++++++++++----------------------- 1 file changed, 40 insertions(+), 54 deletions(-) diff --git a/lib/ecto/repo/schema.ex b/lib/ecto/repo/schema.ex index 086853a3af..a15ce265aa 100644 --- a/lib/ecto/repo/schema.ex +++ b/lib/ecto/repo/schema.ex @@ -175,17 +175,11 @@ defmodule Ecto.Repo.Schema do header = case query.select do - %Ecto.Query.SelectExpr{expr: {:%{}, [], [{:|, _, [{:&, _, [ix]}, args]}]}, fields: fields} -> - {updated_fields, updated_set} = - Enum.map_reduce(args, MapSet.new(), fn {field, _}, set -> - dumped_field = insert_all_select_dump!(field, dumper) - {dumped_field, MapSet.put(set, field)} - end) + %Ecto.Query.SelectExpr{expr: {:%{}, [], [{:|, _, [{:&, _, [ix]}, args]}]}} -> + updated_fields = + Enum.map(args, fn {field, _} -> insert_all_select_dump!(field, dumper) end) - unchanged_fields = - insert_all_source_fields(query, ix, fields, updated_set, length(args), dumper) - - unchanged_fields ++ updated_fields + insert_all_source_fields(query, ix, args, dumper) ++ updated_fields %Ecto.Query.SelectExpr{expr: {:%{}, _ctx, args}} -> Enum.map(args, fn {field, _} -> insert_all_select_dump!(field, dumper) end) @@ -193,8 +187,8 @@ defmodule Ecto.Repo.Schema do %Ecto.Query.SelectExpr{take: %{^ix => {_fun, fields}}} -> Enum.map(fields, &insert_all_select_dump!(&1, dumper)) - %Ecto.Query.SelectExpr{expr: {:&, _, [ix]}, fields: fields} -> - insert_all_source_fields(query, ix, fields, MapSet.new(), 0, dumper) + %Ecto.Query.SelectExpr{expr: {:&, _, [ix]}} -> + insert_all_source_fields(query, ix, [], dumper) _ -> insert_all_select_error!(query) @@ -317,66 +311,58 @@ defmodule Ecto.Repo.Schema do {rows, Enum.reverse(cast_params), counter} end - defp insert_all_source_fields(query, ix, fields, updated_set, updated_count, dumper) do - # Keep insert headers aligned with the planner's SELECT fields by checking the - # physical source column before mapping each logical field to its destination. + defp insert_all_source_fields(query, ix, updates, dumper) do + fields = query.select.fields + count = length(fields) - length(updates) + if count < 0, do: insert_all_select_error!(query) + fields = Enum.take(fields, count) + updates = Map.new(updates) + + # Schema fields are logical names; SELECT fields contain physical column names. case elem(query.sources, ix) do {_, schema, _} when is_atom(schema) and not is_nil(schema) -> - source_fields = + selected = case query.select.take do - %{^ix => {_fun, selected_fields}} -> selected_fields + %{^ix => {_, selected}} -> selected _ -> schema.__schema__(:query_fields) end - |> Enum.filter(&is_atom/1) - |> Enum.reject(&MapSet.member?(updated_set, &1)) - - {source_exprs, updated_exprs} = Enum.split(fields, length(source_fields)) - - if length(source_exprs) != length(source_fields) or length(updated_exprs) != updated_count do - insert_all_select_error!(query) - end source_dumper = schema.__schema__(:dump) - Enum.zip_with(source_fields, source_exprs, fn field, expr -> - source = insert_all_source_field!(query, ix, expr) - {expected_source, _, _} = Map.get(source_dumper, field, {field, :any, :always}) - - if source != expected_source do - insert_all_select_error!(query) + {header, leftover} = + for field <- selected, is_atom(field), not is_map_key(updates, field), reduce: {[], fields} do + {header, fields} -> + source = + case source_dumper do + %{^field => {source, _, _}} -> source + _ -> field + end + + case fields do + [{{:., _, [{:&, _, [^ix]}, ^source]}, [], []} | fields] -> + field = insert_all_select_dump!(if(dumper, do: field, else: source), dumper) + {[field | header], fields} + + _ -> + insert_all_select_error!(query) + end end - if dumper, do: insert_all_select_dump!(field, dumper), else: source - end) + if leftover != [], do: insert_all_select_error!(query) + Enum.reverse(header) _ -> - {source_exprs, updated_exprs} = split_updated_fields(fields, updated_count) - - if length(updated_exprs) != updated_count do - insert_all_select_error!(query) - end - - Enum.map(source_exprs, fn expr -> - field = insert_all_source_field!(query, ix, expr) + Enum.map(fields, fn + {{:., _, [{:&, _, [^ix]}, field]}, [], []} -> + if is_map_key(updates, field), do: insert_all_select_error!(query) + insert_all_select_dump!(field, dumper) - if MapSet.member?(updated_set, field) do + _ -> insert_all_select_error!(query) - end - - insert_all_select_dump!(field, dumper) end) end end - defp split_updated_fields(fields, 0), do: {fields, []} - defp split_updated_fields(fields, count), do: Enum.split(fields, -count) - - defp insert_all_source_field!(query, ix, {{:., _, [{:&, _, [expr_ix]}, field]}, [], []}) do - if expr_ix == ix, do: field, else: insert_all_select_error!(query) - end - - defp insert_all_source_field!(query, _ix, _expr), do: insert_all_select_error!(query) - defp insert_all_select_error!(query) do raise ArgumentError, """ cannot generate a fields list for insert_all from the given source query: