diff options
| author | Alexander Smirnov <[email protected]> | 2025-04-18 00:51:50 +0000 |
|---|---|---|
| committer | Alexander Smirnov <[email protected]> | 2025-04-18 00:51:50 +0000 |
| commit | fcf98cbcba210753db1754ca6e28c295c535ffbb (patch) | |
| tree | 1b529bd303f9c788e4f398933d7deb47cfd8c3b2 | |
| parent | 1f62eb9b72a10d093a2acd13c434a8a94fa695bc (diff) | |
| parent | 8e5325590b3037c576e7f9981903f5112e181ffe (diff) | |
Merge branch 'rightlib' into merge-libs-250418-0050
91 files changed, 1570 insertions, 232 deletions
diff --git a/build/conf/ts/node_modules.conf b/build/conf/ts/node_modules.conf index 29abae4fcca..585da219047 100644 --- a/build/conf/ts/node_modules.conf +++ b/build/conf/ts/node_modules.conf @@ -43,7 +43,7 @@ _PREPARE_DEPS_CMD=$TOUCH_UNIT \ # In case of no deps we need to create empty outputs for graph connectivity _PREPARE_NO_DEPS_CMD=$TOUCH_UNIT \ - && $YMAKE_PYTHON ${input:"build/scripts/touch.py"} \ + && $YMAKE_PYTHON3 ${input:"build/scripts/touch.py"} \ $_PREPARE_DEPS_INOUTS \ ${hide;kv:"pc magenta"} ${hide;kv:"p TS_NODEP"} diff --git a/build/conf/ts/ts.conf b/build/conf/ts/ts.conf index c74e666c261..b6f6fbad7ba 100644 --- a/build/conf/ts/ts.conf +++ b/build/conf/ts/ts.conf @@ -17,7 +17,7 @@ _TS_PROJECT_SETUP_CMD=$EXTRACT_GENTAR TS_CONFIG_PATH=tsconfig.json -EXTRACT_GENTAR=${cwd:BINDIR} $YMAKE_PYTHON ${input:"build/scripts/autotar_gendirs.py"} --unpack --ext .gentar ${ext=.gentar:AUTO_INPUT} +EXTRACT_GENTAR=${cwd:BINDIR} $YMAKE_PYTHON3 ${input:"build/scripts/autotar_gendirs.py"} --unpack --ext .gentar ${ext=.gentar:AUTO_INPUT} ### @usage: TS_CONFIG(ConfigPath) ### diff --git a/build/export_generators/ide-gradle/build.gradle.kts.any.jinja b/build/export_generators/ide-gradle/build.gradle.kts.any.jinja index 2bff0ca3932..79936b6822c 100644 --- a/build/export_generators/ide-gradle/build.gradle.kts.any.jinja +++ b/build/export_generators/ide-gradle/build.gradle.kts.any.jinja @@ -3,16 +3,16 @@ {#- That is why all common macroses here -#} {%- macro PatchRoots(arg, depend = false, output = false) -%} -{#- Always replace (arcadia_root) === (SOURCE_ROOT in ymake) to $project_root in Gradle -#} +{#- Always replace (arcadia_root) === (SOURCE_ROOT in ymake) to $arcadia_root in Gradle -#} {%- if depend -%} -{#- Replace (export_root) === (BUILD_ROOT in ymake) to $project_root in Gradle, because prebuilt tools in arcadia, not in build_root -#} -"{{ arg|replace(export_root, "$project_root")|replace(arcadia_root, "$project_root") }}" +{#- Replace (export_root) === (BUILD_ROOT in ymake) to $arcadia_root in Gradle, because prebuilt tools in arcadia, not in build_root -#} +"{{ arg|replace(export_root, "$arcadia_root")|replace(arcadia_root, "$arcadia_root") }}" {%- elif output and arg[0] != '/' -%} {#- Relative outputs in buildDir -#} "$buildDir/{{ arg }}" {%- else -%} {#- Replace (export_root) === (BUILD_ROOT in ymake) to baseBuildDir in Gradle - root of all build folders for modules -#} -"{{ arg|replace(export_root, "$baseBuildDir")|replace(arcadia_root, "$project_root") }}" +"{{ arg|replace(export_root, "$baseBuildDir")|replace(arcadia_root, "$arcadia_root") }}" {%- endif -%} {%- endmacro -%} @@ -26,13 +26,14 @@ {%- endmacro -%} {%- endif -%} +{%- include "[generator]/common_dir.jinja" -%} + {%- if proto_template -%} {%- include "[generator]/proto_vars.jinja" -%} {%- include "[generator]/proto_import.jinja" -%} {%- include "[generator]/proto_builddir.jinja" -%} {%- include "[generator]/proto_plugins.jinja" -%} {%- include "[generator]/proto_configuration.jinja" -%} -{%- include "[generator]/proto_source_sets.jinja" -%} {%- include "[generator]/protobuf.jinja" -%} {%- include "[generator]/proto_prepare.jinja" -%} {%- include "[generator]/build.gradle.kts.common.jinja" -%} @@ -45,7 +46,6 @@ {%- include "[generator]/kotlin_plugins.jinja" -%} {%- include "[generator]/preview.jinja" -%} {%- include "[generator]/configuration.jinja" -%} -{%- include "[generator]/source_sets.jinja" -%} {%- include "[generator]/test.jinja" -%} {%- include "[generator]/build.gradle.kts.common.jinja" -%} {%- include "[generator]/dependencies.jinja" -%} diff --git a/build/export_generators/ide-gradle/build.gradle.kts.common.jinja b/build/export_generators/ide-gradle/build.gradle.kts.common.jinja index be59a3af157..72e0f4ca0fa 100644 --- a/build/export_generators/ide-gradle/build.gradle.kts.common.jinja +++ b/build/export_generators/ide-gradle/build.gradle.kts.common.jinja @@ -5,6 +5,11 @@ {%- include "[generator]/javac_flags.jinja" -%} {%- include "[generator]/kotlinc_flags.jinja" -%} +{%- include "[generator]/source_sets.jinja" -%} {%- include "[generator]/codegen.jinja" -%} -{%- include "[generator]/javadoc.jinja" -%} +{#- To disable redundant javadoc (it may fail the build) #} + +tasks.withType<Javadoc>().configureEach { + isEnabled = false +} diff --git a/build/export_generators/ide-gradle/builddir.jinja b/build/export_generators/ide-gradle/builddir.jinja index 88a5694d157..9ee3cd071d4 100644 --- a/build/export_generators/ide-gradle/builddir.jinja +++ b/build/export_generators/ide-gradle/builddir.jinja @@ -1,6 +1,6 @@ {#- empty string #} val baseBuildDir = "{{ export_root }}/gradle.build" -buildDir = file(baseBuildDir + "/" + project.path.replace(":", "/")) +buildDir = file(baseBuildDir + "{%- if common_dir %}/{{ common_dir }}{% endif -%}/" + project.path.replace(":", "/")) subprojects { - buildDir = file(baseBuildDir + "/" + project.path.replace(":", "/")) + buildDir = file(baseBuildDir + "{%- if common_dir %}/{{ common_dir }}{% endif -%}/" + project.path.replace(":", "/")) } diff --git a/build/export_generators/ide-gradle/codegen.jinja b/build/export_generators/ide-gradle/codegen.jinja index 3950c9aa166..4df2ad85a59 100644 --- a/build/export_generators/ide-gradle/codegen.jinja +++ b/build/export_generators/ide-gradle/codegen.jinja @@ -2,22 +2,28 @@ {%- if proto_template %} tasks.getByName("prepareMainProtos").dependsOn({{ taskvar }}) -{%- endif %} +{%- endif -%} +{#- Check main target codegen -#} +{%- if varprefix == "codegen" %} tasks.compileJava.configure { dependsOn({{ taskvar }}) } +{%- endif %} tasks.compileTestJava.configure { dependsOn({{ taskvar }}) } -{%- if with_kotlin %} +{%- if with_kotlin -%} +{#- Check main target codegen -#} +{%- if varprefix == "codegen" %} tasks.compileKotlin.configure { dependsOn({{ taskvar }}) } +{%- endif %} tasks.compileTestKotlin.configure { dependsOn({{ taskvar }}) } -{% endif -%} +{% endif -%} {%- endmacro -%} {%- macro ObjDepends(obj) -%} @@ -33,11 +39,13 @@ tasks.getByName("{{ parent_taskvar }}").dependsOn({{ taskvar }}) {%- if target is defined -%} {%- set current_target = target -%} +{#- Main target codegen -#} {%- set varprefix = "codegen" -%} {%- include "[generator]/codegen_current_target.jinja" -%} {%- endif -%} {%- if extra_targets|length -%} {%- for current_target in extra_targets -%} +{#- TestN target codegen -#} {%- set varprefix = "test" + loop.index0|tojson + "Codegen" -%} {%- include "[generator]/codegen_current_target.jinja" -%} {%- endfor -%} diff --git a/build/export_generators/ide-gradle/codegen_run_java_program.jinja b/build/export_generators/ide-gradle/codegen_run_java_program.jinja index e3604f85e62..87fb1bedeaa 100644 --- a/build/export_generators/ide-gradle/codegen_run_java_program.jinja +++ b/build/export_generators/ide-gradle/codegen_run_java_program.jinja @@ -17,7 +17,7 @@ val {{ varprefix }}{{ run['_object_index'] }} = task<JavaExec>("{{ varprefix }}{ {% for classpath in classpaths -%} {%- set rel_file_classpath = classpath|replace('@', '')|replace(export_root, '')|replace(arcadia_root, '') %} - val classpaths = "$project_root/" + File("$project_root{{ rel_file_classpath }}").readText().trim().replace(":", ":$project_root/") + val classpaths = "$arcadia_root/" + File("$arcadia_root{{ rel_file_classpath }}").readText().trim().replace(":", ":$arcadia_root/") classpath = files(classpaths.split(":")) {%- endfor -%} {%- endif %} diff --git a/build/export_generators/ide-gradle/common_dir.jinja b/build/export_generators/ide-gradle/common_dir.jinja new file mode 100644 index 00000000000..a9106fb926e --- /dev/null +++ b/build/export_generators/ide-gradle/common_dir.jinja @@ -0,0 +1,5 @@ +{%- if common_dir -%} +{%- set common_dir_classpath = '":' + common_dir|replace("/", ":") -%} +{%- else -%} +{%- set common_dir_classpath = false -%} +{%- endif -%} diff --git a/build/export_generators/ide-gradle/configuration.jinja b/build/export_generators/ide-gradle/configuration.jinja index b238ca95bfa..c049cba2dc5 100644 --- a/build/export_generators/ide-gradle/configuration.jinja +++ b/build/export_generators/ide-gradle/configuration.jinja @@ -1,5 +1,6 @@ {#- empty string #} -val project_root = "{{ arcadia_root }}" +val arcadia_root = "{{ arcadia_root }}" +val project_root = "{{ project_root }}" {% if mainClass -%} application { diff --git a/build/export_generators/ide-gradle/dependencies.jinja b/build/export_generators/ide-gradle/dependencies.jinja index a869975b766..d5958c36acc 100644 --- a/build/export_generators/ide-gradle/dependencies.jinja +++ b/build/export_generators/ide-gradle/dependencies.jinja @@ -2,11 +2,11 @@ {%- if annotation_processors|length -%} {%- set lomboks = annotation_processors|select('startsWith', 'contrib/java/org/projectlombok/lombok') -%} {%- for lombok in lomboks %} - {{ funcName }}(files("$project_root/{{ lombok }}")) + {{ funcName }}(files("$arcadia_root/{{ lombok }}")) {%- endfor -%} {%- set annotation_processors = annotation_processors|reject('in', lomboks) -%} {%- for annotation_processor in annotation_processors %} - {{ funcName }}(files("$project_root/{{ annotation_processor }}")) + {{ funcName }}(files("$arcadia_root/{{ annotation_processor }}")) {%- endfor -%} {%- endif -%} {%- endmacro -%} @@ -22,15 +22,18 @@ dependencies { {%- endif -%} {%- if library.prebuilt and library.jar and (library.type != "contrib" or build_contribs) %} - implementation(files("$project_root/{{ library.jar }}")) + implementation(files("$arcadia_root/{{ library.jar }}")) {%- else -%} {%- set classpath = library.classpath -%} {%- if classpath|replace('"','') == classpath -%} {%- set classpath = '"' + classpath + '"' -%} {%- endif -%} +{%- include "[generator]/patch_classpath.jinja" -%} {%- if library.type != "contrib" %} -{%- if library.testdep %} - implementation(project(path = ":{{ library.testdep | replace("/", ":") }}", configuration = "testArtifacts")) +{%- if library.testdep -%} +{%- set classpath = '":' + library.testdep | replace("/", ":") + '"' -%} +{%- include "[generator]/patch_classpath.jinja" %} + implementation(project(path = {{ classpath }}, configuration = "testArtifacts")) {%- else %} implementation({{ classpath }}) {%- endif -%} @@ -53,14 +56,17 @@ dependencies { {%- for extra_target in extra_targets -%} {%- for library in extra_target.consumer if library.classpath -%} {%- if library.prebuilt and library.jar and (library.type != "contrib" or build_contribs) %} - testImplementation(files("$project_root/{{ library.jar }}")) + testImplementation(files("$arcadia_root/{{ library.jar }}")) {%- else -%} {%- set classpath = library.classpath -%} {%- if classpath|replace('"','') == classpath -%} {%- set classpath = '"' + classpath + '"' -%} {%- endif %} -{%- if library.type != "contrib" and library.testdep %} - testImplementation(project(path = ":{{ library.testdep | replace("/", ":") }}", configuration = "testArtifacts")) +{%- include "[generator]/patch_classpath.jinja" -%} +{%- if library.type != "contrib" and library.testdep -%} +{%- set classpath = '":' + library.testdep | replace("/", ":") + '"' -%} +{%- include "[generator]/patch_classpath.jinja" %} + testImplementation(project(path = {{ classpath }}, configuration = "testArtifacts")) {%- else %} testImplementation({{ classpath }}) {%- endif -%} diff --git a/build/export_generators/ide-gradle/javadoc.jinja b/build/export_generators/ide-gradle/javadoc.jinja deleted file mode 100644 index 94fc8c750ff..00000000000 --- a/build/export_generators/ide-gradle/javadoc.jinja +++ /dev/null @@ -1,4 +0,0 @@ -{#- To disable redundant javadoc (it may fail the build) #} -tasks.withType<Javadoc>().configureEach { - isEnabled = false -} diff --git a/build/export_generators/ide-gradle/kotlinc_flags.jinja b/build/export_generators/ide-gradle/kotlinc_flags.jinja index fce0eecfcf7..3f6c8d90ad1 100644 --- a/build/export_generators/ide-gradle/kotlinc_flags.jinja +++ b/build/export_generators/ide-gradle/kotlinc_flags.jinja @@ -12,7 +12,7 @@ tasks.withType<KotlinCompile> { compilerOptions { {%- for kotlinc_flag in kotlinc_flags|unique %} - freeCompilerArgs.add("{{ kotlinc_flag|replace(export_root, "$project_root")|replace(arcadia_root, "$project_root") }}") + freeCompilerArgs.add("{{ kotlinc_flag|replace(export_root, "$arcadia_root")|replace(arcadia_root, "$arcadia_root") }}") {%- endfor %} } } diff --git a/build/export_generators/ide-gradle/patch_classpath.jinja b/build/export_generators/ide-gradle/patch_classpath.jinja new file mode 100644 index 00000000000..c653e7b1e06 --- /dev/null +++ b/build/export_generators/ide-gradle/patch_classpath.jinja @@ -0,0 +1,3 @@ +{%- if common_dir_classpath -%} +{%- set classpath = classpath|replace(common_dir_classpath, '"') -%} +{%- endif -%} diff --git a/build/export_generators/ide-gradle/proto_configuration.jinja b/build/export_generators/ide-gradle/proto_configuration.jinja index 5a9554b2aeb..c2cf9d24592 100644 --- a/build/export_generators/ide-gradle/proto_configuration.jinja +++ b/build/export_generators/ide-gradle/proto_configuration.jinja @@ -1,5 +1,6 @@ {#- empty string #} -val project_root = "{{ arcadia_root }}" +val arcadia_root = "{{ arcadia_root }}" +val project_root = "{{ project_root }}" java { withSourcesJar() diff --git a/build/export_generators/ide-gradle/proto_dependencies.jinja b/build/export_generators/ide-gradle/proto_dependencies.jinja index 61bcc05fa98..dc3ca68e4be 100644 --- a/build/export_generators/ide-gradle/proto_dependencies.jinja +++ b/build/export_generators/ide-gradle/proto_dependencies.jinja @@ -2,12 +2,13 @@ dependencies { {%- for library in target.consumer if library.classpath -%} {%- if library.prebuilt and library.jar and (library.type != "contrib" or target.handler.build_contribs) %} - implementation(files("$project_root/{{ library.jar }}")) + implementation(files("$arcadia_root/{{ library.jar }}")) {%- else -%} {%- set classpath = library.classpath -%} {%- if classpath|replace('"','') == classpath -%} {%- set classpath = '"' + classpath + '"' -%} {%- endif %} +{%- include "[generator]/patch_classpath.jinja" -%} {%- if library.type != "contrib" %} implementation {%- else %} diff --git a/build/export_generators/ide-gradle/proto_prepare.jinja b/build/export_generators/ide-gradle/proto_prepare.jinja index 804428a964e..7be93a58c35 100644 --- a/build/export_generators/ide-gradle/proto_prepare.jinja +++ b/build/export_generators/ide-gradle/proto_prepare.jinja @@ -2,7 +2,7 @@ val prepareMainProtos = tasks.register<Copy>("prepareMainProtos") { {%- if target.proto_files|length %} - from("$project_root") { + from("$arcadia_root") { {#- list of all current project proto files -#} {%- for proto in target.proto_files %} include("{{ proto }}") @@ -30,11 +30,11 @@ val prepareMainProtos = tasks.register<Copy>("prepareMainProtos") { } {%- endif -%} -{%- if extractLibrariesProtosTask -%} +{%- if extractLibrariesProtosTask %} val extractMainLibrariesProtos = tasks.register<Copy>("extractMainLibrariesProtos") { -{%- if libraries|length -%} - from("$project_root") { +{%- if libraries|length %} + from("$arcadia_root") { {#- list of all library directories -#} {%- for library in libraries -%} {%- set path_and_jar = rsplit(library.jar, '/', 2) %} diff --git a/build/export_generators/ide-gradle/proto_source_sets.jinja b/build/export_generators/ide-gradle/proto_source_sets.jinja deleted file mode 100644 index 540b4003295..00000000000 --- a/build/export_generators/ide-gradle/proto_source_sets.jinja +++ /dev/null @@ -1,37 +0,0 @@ -{#- empty string #} -sourceSets { - main { -{%- if target.jar_source_set|length -%} -{%- for source_set in target.jar_source_set -%} -{%- set srcdir_glob = split(source_set, ':') -%} -{%- set srcdir = srcdir_glob[0] -%} -{%- if srcdir != 'src/main/java' %} - java.srcDir({{ PatchRoots(srcdir) }}) -{%- endif -%} -{%- endfor -%} -{%- endif %} -{%- if target.jar_resource_set|length -%} -{%- for resource_set in target.jar_resource_set -%} -{%- set resdir_glob = split(resource_set, ':') -%} -{%- set resdir = resdir_glob[0] -%} -{%- if resdir != 'src/main/resources' %} - resources.srcDir({{ PatchRoots(resdir) }}) -{%- endif -%} -{%- endfor -%} -{%- endif %} - java.srcDir("$buildDir/generated/source/proto/main/java") -{%- if target.proto_grpc %} - java.srcDir("$buildDir/generated/source/proto/main/grpc") -{%- endif %} - } - test { - java.srcDir("$buildDir/generated/source/proto/test/java") -{%- if target.proto_grpc %} - java.srcDir("$buildDir/generated/source/proto/test/grpc") -{%- endif %} - } -} - -tasks.withType<Jar>() { - duplicatesStrategy = DuplicatesStrategy.INCLUDE -} diff --git a/build/export_generators/ide-gradle/settings.gradle.kts.jinja b/build/export_generators/ide-gradle/settings.gradle.kts.jinja index 68bb5b594d9..865019f9017 100644 --- a/build/export_generators/ide-gradle/settings.gradle.kts.jinja +++ b/build/export_generators/ide-gradle/settings.gradle.kts.jinja @@ -1,12 +1,15 @@ +{%- include "[generator]/common_dir.jinja" -%} rootProject.name = "{{ project_name }}" - {% for subdir in subdirs -%} {%- set arcadia_subdir = arcadia_root + "/" + subdir -%} {%- if arcadia_subdir != project_root -%} -{%- set classname = subdir | replace("/", ":") %} -include(":{{ classname }}") -project(":{{ classname }}").projectDir = file("{{ arcadia_subdir }}") -{% endif -%} +{%- set classpath = '":' + subdir | replace("/", ":") + '"' -%} +{%- include "[generator]/patch_classpath.jinja" %} +include({{ classpath }}) +{%- if not common_dir_classpath %} +project({{ classpath }}).projectDir = file("{{ arcadia_subdir }}") +{%- endif -%} +{%- endif -%} {%- endfor -%} {%- include "[generator]/debug.jinja" ignore missing -%} diff --git a/build/export_generators/ide-gradle/source_sets.jinja b/build/export_generators/ide-gradle/source_sets.jinja index aff24144a88..e127569c458 100644 --- a/build/export_generators/ide-gradle/source_sets.jinja +++ b/build/export_generators/ide-gradle/source_sets.jinja @@ -1,6 +1,8 @@ {#- empty string #} sourceSets { -{%- if target.runs|length or target.jar_source_set|length %} +{%- set target_jar_source_set = target.jar_source_set|reject('startsWith', 'src/main/java:')|unique -%} +{%- set target_jar_resource_set = target.jar_resource_set|reject('startsWith', 'src/main/resources:')|unique -%} +{%- if proto_template or target_jar_source_set|length or target_jar_resource_set|length %} main { {#- Default by Gradle: @@ -9,23 +11,25 @@ sourceSets { resources.srcDir("src/main/resources") #} -{%- if target.jar_source_set|length -%} -{%- for source_set in target.jar_source_set -%} +{%- if target_jar_source_set|length -%} +{%- for source_set in target_jar_source_set -%} {%- set srcdir_glob = split(source_set, ':') -%} -{%- set srcdir = srcdir_glob[0] -%} -{%- if srcdir != 'src/main/java' %} +{%- set srcdir = srcdir_glob[0] %} java.srcDir({{ PatchRoots(srcdir) }}) -{%- endif -%} {%- endfor -%} {%- endif %} -{%- if target.jar_resource_set|length -%} -{%- for resource_set in target.jar_resource_set -%} +{%- if target_jar_resource_set|length -%} +{%- for resource_set in target_jar_resource_set -%} {%- set resdir_glob = split(resource_set, ':') -%} -{%- set resdir = resdir_glob[0] -%} -{%- if resdir != 'src/main/resources' %} +{%- set resdir = resdir_glob[0] %} resources.srcDir({{ PatchRoots(resdir) }}) -{%- endif -%} {%- endfor -%} +{%- endif -%} +{%- if proto_template %} + java.srcDir("$buildDir/generated/source/proto/main/java") +{%- if target.proto_grpc %} + java.srcDir("$buildDir/generated/source/proto/main/grpc") +{%- endif %} {%- endif %} } {%- endif %} @@ -37,6 +41,12 @@ sourceSets { resources.srcDir("src/test/resources") #} +{%- if proto_template %} + java.srcDir("$buildDir/generated/source/proto/test/java") +{%- if target.proto_grpc %} + java.srcDir("$buildDir/generated/source/proto/test/grpc") +{%- endif -%} +{%- else %} java.srcDir("ut/java") resources.srcDir("ut/resources") java.srcDir("src/test-integration/java") @@ -48,25 +58,22 @@ sourceSets { java.srcDir("src/intTest/java") resources.srcDir("src/intTest/resources") -{%- set extra_target_source_sets = extra_targets|selectattr('jar_source_set')|map(attribute='jar_source_set')|sum|unique -%} -{%- if extra_target_source_sets|length -%} -{%- for source_set in extra_target_source_sets -%} -{%- set srcdir_glob = split(source_set, ':') -%} -{%- set srcdir = srcdir_glob[0] -%} -{%- if srcdir != 'src/test/java' %} +{%- set extra_target_source_sets = extra_targets|selectattr('jar_source_set')|map(attribute='jar_source_set')|sum|reject('startsWith', 'src/test/java:')|unique -%} +{%- if extra_target_source_sets|length -%} +{%- for source_set in extra_target_source_sets -%} +{%- set srcdir_glob = split(source_set, ':') -%} +{%- set srcdir = srcdir_glob[0] %} java.srcDir({{ PatchRoots(srcdir) }}) -{%- endif -%} -{%- endfor -%} -{%- endif %} -{%- set extra_target_resource_sets = extra_targets|selectattr('jar_resource_set')|map(attribute='jar_resource_set')|sum|unique -%} -{%- if extra_target_resource_sets|length -%} -{%- for resource_set in extra_target_resource_sets -%} -{%- set resdir_glob = split(resource_set, ':') -%} -{%- set resdir = resdir_glob[0] -%} -{%- if resdir != 'src/main/resources' %} +{%- endfor -%} +{%- endif %} +{%- set extra_target_resource_sets = extra_targets|selectattr('jar_resource_set')|map(attribute='jar_resource_set')|sum|reject('startsWith', 'src/test/resources:')|unique -%} +{%- if extra_target_resource_sets|length -%} +{%- for resource_set in extra_target_resource_sets -%} +{%- set resdir_glob = split(resource_set, ':') -%} +{%- set resdir = resdir_glob[0] %} resources.srcDir({{ PatchRoots(resdir) }}) -{%- endif -%} -{%- endfor -%} +{%- endfor -%} +{%- endif -%} {%- endif %} } } diff --git a/build/plugins/_dart_fields.py b/build/plugins/_dart_fields.py index 64232696d15..37db4c5a435 100644 --- a/build/plugins/_dart_fields.py +++ b/build/plugins/_dart_fields.py @@ -1116,6 +1116,7 @@ class TestFiles: 'maps/renderer/libs/data_sets/yt_data_set', 'maps/renderer/libs/design', 'maps/renderer/libs/geosx', + 'maps/renderer/libs/geojson_to_yt', 'maps/renderer/libs/gltf', 'maps/renderer/libs/golden', 'maps/renderer/libs/hd3d', diff --git a/build/ymake.core.conf b/build/ymake.core.conf index 8adf8bc5141..18ad281440f 100644 --- a/build/ymake.core.conf +++ b/build/ymake.core.conf @@ -1692,6 +1692,7 @@ module EXECTEST: _BARE_UNIT { SET(MODULE_SUFFIX .pkg.fake) SETUP_EXECTEST() SET_APPEND(_MAKEFILE_INCLUDE_LIKE_DEPS canondata/result.json) + DISABLE(_NEED_SBOM_INFO) } # tag:cpp-specific tag:test @@ -2352,6 +2353,7 @@ multimodule PACKAGE { } SET(NEED_PLATFORM_PEERDIRS no) SET(_COPY_FILE_CONTEXT TEXT) + DISABLE(_NEED_SBOM_INFO) } module PACKAGE_UNION: UNION { .CMD=UNION_CMD @@ -2410,6 +2412,7 @@ module UNION: _BASE_UNIT { SET(NEED_PLATFORM_PEERDIRS no) PEERDIR_TAGS=CPP_PROTO CPP_PROTO_FROM_SCHEMA CPP_FBS PY2 PY2_NATIVE PY3_NATIVE YQL_UDF_SHARED __EMPTY__ RESOURCE_LIB DOCSBOOK JAR_RUNNABLE PY3_BIN DLL TS PACKAGE_UNION + DISABLE(_NEED_SBOM_INFO) UNION_OUTS=${hide;late_out:AUTO_INPUT} when ($_UNION_EXPLICIT_OUTPUTS) { UNION_OUTS=$_EXPAND_INS_OUTS($_UNION_EXPLICIT_OUTPUTS) diff --git a/library/cpp/containers/dense_hash/dense_hash.h b/library/cpp/containers/dense_hash/dense_hash.h index 739479c25a3..b5feb16eefb 100644 --- a/library/cpp/containers/dense_hash/dense_hash.h +++ b/library/cpp/containers/dense_hash/dense_hash.h @@ -168,14 +168,10 @@ public: } else { initSize = FastClp2(initSize); } - BucketMask = initSize - 1; + Buckets.clear(); + BucketMask = 0; NumFilled = 0; - TVector<value_type> tmp; - for (size_type i = 0; i < initSize; ++i) { - tmp.emplace_back(EmptyMarker, mapped_type{}); - } - tmp.swap(Buckets); - GrowThreshold = Max<size_type>(1, initSize * MaxLoadFactor / 100) - 1; + Grow(initSize); } template <class K> diff --git a/library/cpp/monlib/dynamic_counters/page.cpp b/library/cpp/monlib/dynamic_counters/page.cpp index 5cd750026fb..73b1309d814 100644 --- a/library/cpp/monlib/dynamic_counters/page.cpp +++ b/library/cpp/monlib/dynamic_counters/page.cpp @@ -4,6 +4,7 @@ #include <library/cpp/monlib/service/pages/templates.h> #include <library/cpp/string_utils/quote/quote.h> +#include <util/string/builder.h> #include <util/string/split.h> #include <util/system/tls.h> @@ -26,6 +27,19 @@ TMaybe<EFormat> ParseFormat(TStringBuf str) { } } +namespace { + +TStringBuf GetParams(NMonitoring::IMonHttpRequest& request) { + TStringBuf uri = request.GetUri(); + TStringBuf params = uri.After('?'); + if (params.Size() == uri.Size()) { + params.Clear(); + } + return params; +} + +} + void TDynamicCountersPage::Output(NMonitoring::IMonHttpRequest& request) { if (OutputCallback) { OutputCallback(); @@ -37,28 +51,51 @@ void TDynamicCountersPage::Output(NMonitoring::IMonHttpRequest& request) { }; TVector<TStringBuf> parts; - StringSplitter(request.GetPathInfo()) - .Split('/') - .SkipEmpty() - .Collect(&parts); + TMaybe<EFormat> format; + TStringBuf params = GetParams(request); - TMaybe<EFormat> format = !parts.empty() ? ParseFormat(parts.back()) : Nothing(); - if (format) { - parts.pop_back(); - } + if (request.GetPathInfo().empty() && !params.empty()) { + StringSplitter(params).Split('&').SkipEmpty().Consume([&](TStringBuf part) { + TStringBuf name; + TStringBuf value; + part.Split('=', name, value); + if (name.StartsWith("@")) { + if (name == "@format") { + format = ParseFormat(value); + } else if (name == "@name_label") { + nameLabel = value; + } else if (name == "@private") { + visibility = TCountableBase::EVisibility::Private; + } + } else { + parts.push_back(part); + } + return true; + }); + } else { + StringSplitter(request.GetPathInfo()) + .Split('/') + .SkipEmpty() + .Collect(&parts); - if (!parts.empty() && parts.back().StartsWith(TStringBuf("name_label="))) { - TVector<TString> labels; - StringSplitter(parts.back()).Split('=').SkipEmpty().Collect(&labels); - if (labels.size() == 2U) { - nameLabel = labels.back(); + format = !parts.empty() ? ParseFormat(parts.back()) : Nothing(); + if (format) { + parts.pop_back(); } - parts.pop_back(); - } - if (!parts.empty() && parts.back() == TStringBuf("private")) { - visibility = TCountableBase::EVisibility::Private; - parts.pop_back(); + if (!parts.empty() && parts.back().StartsWith(TStringBuf("name_label="))) { + TVector<TString> labels; + StringSplitter(parts.back()).Split('=').SkipEmpty().Collect(&labels); + if (labels.size() == 2U) { + nameLabel = labels.back(); + } + parts.pop_back(); + } + + if (!parts.empty() && parts.back() == TStringBuf("private")) { + visibility = TCountableBase::EVisibility::Private; + parts.pop_back(); + } } auto counters = Counters; @@ -121,9 +158,15 @@ void TDynamicCountersPage::HandleAbsentSubgroup(IMonHttpRequest& request) { void TDynamicCountersPage::BeforePre(IMonHttpRequest& request) { IOutputStream& out = request.Output(); + TStringBuf params = GetParams(request); + TStringBuilder base; + base << Path << '?'; + if (!params.empty()) { + base << params << '&'; + } HTML(out) { DIV() { - out << "<a href='" << request.GetPath() << "/json'>Counters as JSON</a>"; + out << "<a href='" << base << "@format=json'>Counters as JSON</a>"; out << " for Solomon"; } @@ -133,9 +176,11 @@ void TDynamicCountersPage::BeforePre(IMonHttpRequest& request) { UL() { currentCounters->EnumerateSubgroups([&](const TString& name, const TString& value) { LI() { - TString pathPart = name + "=" + value; - Quote(pathPart, ""); - out << "\n<a href='" << request.GetPath() << "/" << pathPart << "'>" << name << " " << value << "</a>"; + auto escName = name; + auto escValue = value; + Quote(escName); + Quote(escValue); + out << "\n<a href='" << base << escName << '=' << escValue << "'>" << name << " " << value << "</a>"; } }); } diff --git a/library/cpp/tld/tlds-alpha-by-domain.txt b/library/cpp/tld/tlds-alpha-by-domain.txt index 3353c0add62..88d9cd720d3 100644 --- a/library/cpp/tld/tlds-alpha-by-domain.txt +++ b/library/cpp/tld/tlds-alpha-by-domain.txt @@ -1,4 +1,4 @@ -# Version 2025041300, Last Updated Sun Apr 13 07:07:01 2025 UTC +# Version 2025041502, Last Updated Wed Apr 16 07:07:01 2025 UTC AAA AARP ABB diff --git a/library/cpp/yt/threading/atomic_object.h b/library/cpp/yt/threading/atomic_object.h index 8b642c0f4fb..a77ade0a00d 100644 --- a/library/cpp/yt/threading/atomic_object.h +++ b/library/cpp/yt/threading/atomic_object.h @@ -1,6 +1,6 @@ #pragma once -#include <library/cpp/yt/threading/rw_spin_lock.h> +#include <library/cpp/yt/threading/writer_starving_rw_spin_lock.h> #include <concepts> diff --git a/library/cpp/yt/threading/rw_spin_lock-inl.h b/library/cpp/yt/threading/rw_spin_lock-inl.h index 779de1b64a8..0a31b1d9dec 100644 --- a/library/cpp/yt/threading/rw_spin_lock-inl.h +++ b/library/cpp/yt/threading/rw_spin_lock-inl.h @@ -31,7 +31,7 @@ inline void TReaderWriterSpinLock::AcquireReaderForkFriendly() noexcept inline void TReaderWriterSpinLock::ReleaseReader() noexcept { auto prevValue = Value_.fetch_sub(ReaderDelta, std::memory_order::release); - Y_ASSERT((prevValue & ~WriterMask) != 0); + Y_ASSERT((prevValue & ~(WriterMask | WriterReadyMask)) != 0); NDetail::RecordSpinLockReleased(); } @@ -45,14 +45,14 @@ inline void TReaderWriterSpinLock::AcquireWriter() noexcept inline void TReaderWriterSpinLock::ReleaseWriter() noexcept { - auto prevValue = Value_.fetch_and(~WriterMask, std::memory_order::release); + auto prevValue = Value_.fetch_and(~(WriterMask | WriterReadyMask), std::memory_order::release); Y_ASSERT(prevValue & WriterMask); NDetail::RecordSpinLockReleased(); } inline bool TReaderWriterSpinLock::IsLocked() const noexcept { - return Value_.load() != UnlockedValue; + return (Value_.load() & ~WriterReadyMask) != 0; } inline bool TReaderWriterSpinLock::IsLockedByReader() const noexcept @@ -68,7 +68,7 @@ inline bool TReaderWriterSpinLock::IsLockedByWriter() const noexcept inline bool TReaderWriterSpinLock::TryAcquireReader() noexcept { auto oldValue = Value_.fetch_add(ReaderDelta, std::memory_order::acquire); - if ((oldValue & WriterMask) != 0) { + if ((oldValue & (WriterMask | WriterReadyMask)) != 0) { Value_.fetch_sub(ReaderDelta, std::memory_order::relaxed); return false; } @@ -79,7 +79,7 @@ inline bool TReaderWriterSpinLock::TryAcquireReader() noexcept inline bool TReaderWriterSpinLock::TryAndTryAcquireReader() noexcept { auto oldValue = Value_.load(std::memory_order::relaxed); - if ((oldValue & WriterMask) != 0) { + if ((oldValue & (WriterMask | WriterReadyMask)) != 0) { return false; } return TryAcquireReader(); @@ -88,7 +88,7 @@ inline bool TReaderWriterSpinLock::TryAndTryAcquireReader() noexcept inline bool TReaderWriterSpinLock::TryAcquireReaderForkFriendly() noexcept { auto oldValue = Value_.load(std::memory_order::relaxed); - if ((oldValue & WriterMask) != 0) { + if ((oldValue & (WriterMask | WriterReadyMask)) != 0) { return false; } auto newValue = oldValue + ReaderDelta; @@ -98,22 +98,35 @@ inline bool TReaderWriterSpinLock::TryAcquireReaderForkFriendly() noexcept return acquired; } -inline bool TReaderWriterSpinLock::TryAcquireWriter() noexcept +inline bool TReaderWriterSpinLock::TryAcquireWriterWithExpectedValue(TValue expected) noexcept { - auto expected = UnlockedValue; - - bool acquired = Value_.compare_exchange_weak(expected, WriterMask, std::memory_order::acquire); + bool acquired = Value_.compare_exchange_weak(expected, WriterMask, std::memory_order::acquire); NDetail::RecordSpinLockAcquired(acquired); return acquired; } +inline bool TReaderWriterSpinLock::TryAcquireWriter() noexcept +{ + // NB(pavook): we cannot expect writer ready to be set, as this method + // might be called without indicating writer readiness and we cannot + // indicate readiness on the hot path. This means that code calling + // TryAcquireWriter will spin against code calling AcquireWriter. + return TryAcquireWriterWithExpectedValue(UnlockedValue); +} + inline bool TReaderWriterSpinLock::TryAndTryAcquireWriter() noexcept { auto oldValue = Value_.load(std::memory_order::relaxed); - if (oldValue != UnlockedValue) { + + if ((oldValue & WriterReadyMask) == 0) { + oldValue = Value_.fetch_or(WriterReadyMask, std::memory_order::relaxed); + } + + if ((oldValue & (~WriterReadyMask)) != 0) { return false; } - return TryAcquireWriter(); + + return TryAcquireWriterWithExpectedValue(WriterReadyMask); } //////////////////////////////////////////////////////////////////////////////// diff --git a/library/cpp/yt/threading/rw_spin_lock.h b/library/cpp/yt/threading/rw_spin_lock.h index a915e677e82..64a241bb6b4 100644 --- a/library/cpp/yt/threading/rw_spin_lock.h +++ b/library/cpp/yt/threading/rw_spin_lock.h @@ -16,8 +16,23 @@ namespace NYT::NThreading { //! Single-writer multiple-readers spin lock. /*! - * Reader-side calls are pretty cheap. - * The lock is unfair. + * Reader-side acquires are pretty cheap, and readers don't spin unless writers + * are present. + * + * The lock is unfair, but writers are prioritized over readers, that is, + * if AcquireWriter() is called at some time, then some writer + * (not necessarily the same one that called AcquireWriter) will succeed + * in the next time. This is implemented by an additional flag "WriterReady", + * that writers set on arrival. No readers can proceed until this flag is reset. + * + * WARNING: You probably should not use this lock if forks are possible: see + * fork_aware_rw_spin_lock.h for a proper fork-safe lock which does the housekeeping for you. + * + * WARNING: This lock is not recursive: you can't call AcquireReader() twice in the same + * thread, as that may lead to a deadlock. For the same reason you shouldn't do WaitFor or any + * other context switch under lock. + * + * See tla+/spinlock.tla for the formally verified lock's properties. */ class TReaderWriterSpinLock : public TSpinLockBase @@ -29,18 +44,26 @@ public: /*! * Optimized for the case of read-intensive workloads. * Cheap (just one atomic increment and no spinning if no writers are present). - * Don't use this call if forks are possible: forking at some + * + * WARNING: Don't use this call if forks are possible: forking at some * intermediate point inside #AcquireReader may corrupt the lock state and - * leave lock forever stuck for the child process. + * leave the lock stuck forever for the child process. + * + * WARNING: The lock is not recursive/reentrant, i.e. it assumes that no thread calls + * AcquireReader() if the reader is already acquired for it. */ void AcquireReader() noexcept; //! Acquires the reader lock. /*! * A more expensive version of #AcquireReader (includes at least * one atomic load and CAS; also may spin even if just readers are present). + * * In contrast to #AcquireReader, this method can be used in the presence of forks. - * Note that fork-friendliness alone does not provide fork-safety: additional - * actions must be performed to release the lock after a fork. + * + * WARNING: fork-friendliness alone does not provide fork-safety: additional + * actions must be performed to release the lock after a fork. This means you + * probably should NOT use this lock in the presence of forks, consider + * fork_aware_rw_spin_lock.h instead as a proper fork-safe lock. */ void AcquireReaderForkFriendly() noexcept; //! Tries acquiring the reader lock; see #AcquireReader. @@ -94,10 +117,12 @@ private: using TValue = ui32; static constexpr TValue UnlockedValue = 0; static constexpr TValue WriterMask = 1; - static constexpr TValue ReaderDelta = 2; + static constexpr TValue WriterReadyMask = 2; + static constexpr TValue ReaderDelta = 4; std::atomic<TValue> Value_ = UnlockedValue; + bool TryAcquireWriterWithExpectedValue(TValue expected) noexcept; bool TryAndTryAcquireReader() noexcept; bool TryAndTryAcquireWriter() noexcept; diff --git a/library/cpp/yt/threading/unittests/rw_spin_lock_ut.cpp b/library/cpp/yt/threading/unittests/rw_spin_lock_ut.cpp new file mode 100644 index 00000000000..653772604ce --- /dev/null +++ b/library/cpp/yt/threading/unittests/rw_spin_lock_ut.cpp @@ -0,0 +1,56 @@ +#include <library/cpp/testing/gtest/gtest.h> + +#include <library/cpp/yt/threading/rw_spin_lock.h> + +#include <util/thread/pool.h> + +#include <latch> +#include <thread> + +namespace NYT::NThreading { +namespace { + +//////////////////////////////////////////////////////////////////////////////// + +TEST(TReaderWriterSpinLockTest, WriterPriority) +{ + int readerThreads = 10; + std::latch latch(readerThreads + 1); + std::atomic<size_t> finishedCount = {0}; + + TReaderWriterSpinLock lock; + + volatile std::atomic<ui32> x = {0}; + + auto readerTask = [&latch, &lock, &finishedCount, &x] () { + latch.arrive_and_wait(); + while (true) { + { + auto guard = ReaderGuard(lock); + // do some stuff + for (ui32 i = 0; i < 10'000u; ++i) { + x.fetch_add(i); + } + } + if (finishedCount.fetch_add(1) > 20'000) { + break; + } + } + }; + + auto readerPool = CreateThreadPool(readerThreads); + for (int i = 0; i < readerThreads; ++i) { + readerPool->SafeAddFunc(readerTask); + } + + latch.arrive_and_wait(); + while (finishedCount.load() == 0); + auto guard = WriterGuard(lock); + EXPECT_LE(finishedCount.load(), 1'000u); + DoNotOptimizeAway(x); +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace +} // namespace NYT::NConcurrency diff --git a/library/cpp/yt/threading/unittests/spin_lock_fork_ut.cpp b/library/cpp/yt/threading/unittests/spin_lock_fork_ut.cpp new file mode 100644 index 00000000000..26e58fff745 --- /dev/null +++ b/library/cpp/yt/threading/unittests/spin_lock_fork_ut.cpp @@ -0,0 +1,160 @@ +#include <library/cpp/testing/gtest/gtest.h> + +#include <library/cpp/yt/threading/rw_spin_lock.h> +#include <library/cpp/yt/threading/fork_aware_spin_lock.h> + +#include <util/thread/pool.h> + +#include <sys/wait.h> + +namespace NYT::NThreading { +namespace { + +//////////////////////////////////////////////////////////////////////////////// + +TEST(TReaderWriterSpinLockTest, ForkFriendlyness) +{ + std::atomic<bool> stopped = {false}; + YT_DECLARE_SPIN_LOCK(TReaderWriterSpinLock, lock); + + auto readerTask = [&lock, &stopped] () { + while (!stopped.load()) { + ForkFriendlyReaderGuard(lock); + } + }; + + auto tryReaderTask = [&lock, &stopped] () { + while (!stopped.load()) { + // NB(pavook): TryAcquire instead of Acquire to minimize checks. + bool acquired = lock.TryAcquireReaderForkFriendly(); + if (acquired) { + lock.ReleaseReader(); + } + } + }; + + auto tryWriterTask = [&lock, &stopped] () { + while (!stopped.load()) { + Sleep(TDuration::MicroSeconds(1)); + bool acquired = lock.TryAcquireWriter(); + if (acquired) { + lock.ReleaseWriter(); + } + } + }; + + auto writerTask = [&lock, &stopped] () { + while (!stopped.load()) { + Sleep(TDuration::MicroSeconds(1)); + WriterGuard(lock); + } + }; + + int readerCount = 20; + int writerCount = 10; + + auto reader = CreateThreadPool(readerCount); + auto writer = CreateThreadPool(writerCount); + + for (int i = 0; i < readerCount / 2; ++i) { + reader->SafeAddFunc(readerTask); + reader->SafeAddFunc(tryReaderTask); + } + for (int i = 0; i < writerCount / 2; ++i) { + writer->SafeAddFunc(writerTask); + writer->SafeAddFunc(tryWriterTask); + } + + // And let the chaos begin! + int forkCount = 2000; + for (int iter = 1; iter <= forkCount; ++iter) { + pid_t pid; + { + auto guard = WriterGuard(lock); + pid = fork(); + } + + YT_VERIFY(pid >= 0); + + // NB(pavook): check different orders to maximize chaos. + if (iter % 2 == 0) { + ReaderGuard(lock); + } + WriterGuard(lock); + ReaderGuard(lock); + if (pid == 0) { + // NB(pavook): thread pools are no longer with us. + _exit(0); + } + } + + for (int i = 1; i <= forkCount; ++i) { + int status; + YT_VERIFY(waitpid(0, &status, 0) > 0); + YT_VERIFY(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } + + stopped.store(true); +} + +//////////////////////////////////////////////////////////////////////////////// + +TEST(TForkAwareSpinLockTest, ForkSafety) +{ + std::atomic<bool> stopped = {false}; + YT_DECLARE_SPIN_LOCK(TForkAwareSpinLock, lock); + + auto acquireTask = [&lock, &stopped] () { + while (!stopped.load()) { + Guard(lock); + } + }; + + // NB(pavook): TryAcquire instead of Acquire to minimize checks. + auto tryAcquireTask = [&lock, &stopped] () { + while (!stopped.load()) { + bool acquired = lock.TryAcquire(); + if (acquired) { + lock.Release(); + } + } + }; + + int workerCount = 20; + + auto worker = CreateThreadPool(workerCount); + + for (int i = 0; i < workerCount / 2; ++i) { + worker->SafeAddFunc(acquireTask); + worker->SafeAddFunc(tryAcquireTask); + } + + // And let the chaos begin! + int forkCount = 2000; + for (int iter = 1; iter <= forkCount; ++iter) { + pid_t pid = fork(); + + YT_VERIFY(pid >= 0); + + Guard(lock); + Guard(lock); + + if (pid == 0) { + // NB(pavook): thread pools are no longer with us. + _exit(0); + } + } + + for (int i = 1; i <= forkCount; ++i) { + int status; + YT_VERIFY(waitpid(0, &status, 0) > 0); + YT_VERIFY(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } + + stopped.store(true); +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace +} // namespace NYT::NConcurrency diff --git a/library/cpp/yt/threading/unittests/ya.make b/library/cpp/yt/threading/unittests/ya.make index ef9b5d29951..da006012c00 100644 --- a/library/cpp/yt/threading/unittests/ya.make +++ b/library/cpp/yt/threading/unittests/ya.make @@ -5,9 +5,14 @@ INCLUDE(${ARCADIA_ROOT}/library/cpp/yt/ya_cpp.make.inc) SRCS( count_down_latch_ut.cpp recursive_spin_lock_ut.cpp + rw_spin_lock_ut.cpp spin_wait_ut.cpp ) +IF (NOT OS_WINDOWS) + SRC(spin_lock_fork_ut.cpp) +ENDIF() + PEERDIR( library/cpp/yt/assert library/cpp/yt/threading diff --git a/library/cpp/yt/threading/writer_starving_rw_spin_lock-inl.h b/library/cpp/yt/threading/writer_starving_rw_spin_lock-inl.h new file mode 100644 index 00000000000..cf8bde715cc --- /dev/null +++ b/library/cpp/yt/threading/writer_starving_rw_spin_lock-inl.h @@ -0,0 +1,101 @@ +#pragma once +#ifndef WRITER_STARVING_RW_SPIN_LOCK_INL_H_ +#error "Direct inclusion of this file is not allowed, include rw_spin_lock.h" +// For the sake of sane code completion. +#include "writer_starving_rw_spin_lock.h" +#endif +#undef WRITER_STARVING_RW_SPIN_LOCK_INL_H_ + +#include "spin_wait.h" + +namespace NYT::NThreading { + +//////////////////////////////////////////////////////////////////////////////// + +inline void TWriterStarvingRWSpinLock::AcquireReader() noexcept +{ + if (TryAcquireReader()) { + return; + } + AcquireReaderSlow(); +} + +inline void TWriterStarvingRWSpinLock::ReleaseReader() noexcept +{ + auto prevValue = Value_.fetch_sub(ReaderDelta, std::memory_order::release); + Y_ASSERT((prevValue & ~WriterMask) != 0); + NDetail::RecordSpinLockReleased(); +} + +inline void TWriterStarvingRWSpinLock::AcquireWriter() noexcept +{ + if (TryAcquireWriter()) { + return; + } + AcquireWriterSlow(); +} + +inline void TWriterStarvingRWSpinLock::ReleaseWriter() noexcept +{ + auto prevValue = Value_.fetch_and(~WriterMask, std::memory_order::release); + Y_ASSERT(prevValue & WriterMask); + NDetail::RecordSpinLockReleased(); +} + +inline bool TWriterStarvingRWSpinLock::IsLocked() const noexcept +{ + return Value_.load() != UnlockedValue; +} + +inline bool TWriterStarvingRWSpinLock::IsLockedByReader() const noexcept +{ + return Value_.load() >= ReaderDelta; +} + +inline bool TWriterStarvingRWSpinLock::IsLockedByWriter() const noexcept +{ + return (Value_.load() & WriterMask) != 0; +} + +inline bool TWriterStarvingRWSpinLock::TryAcquireReader() noexcept +{ + auto oldValue = Value_.fetch_add(ReaderDelta, std::memory_order::acquire); + if ((oldValue & WriterMask) != 0) { + Value_.fetch_sub(ReaderDelta, std::memory_order::relaxed); + return false; + } + NDetail::RecordSpinLockAcquired(); + return true; +} + +inline bool TWriterStarvingRWSpinLock::TryAndTryAcquireReader() noexcept +{ + auto oldValue = Value_.load(std::memory_order::relaxed); + if ((oldValue & WriterMask) != 0) { + return false; + } + return TryAcquireReader(); +} + +inline bool TWriterStarvingRWSpinLock::TryAcquireWriter() noexcept +{ + auto expected = UnlockedValue; + + bool acquired = Value_.compare_exchange_weak(expected, WriterMask, std::memory_order::acquire); + NDetail::RecordSpinLockAcquired(acquired); + return acquired; +} + +inline bool TWriterStarvingRWSpinLock::TryAndTryAcquireWriter() noexcept +{ + auto oldValue = Value_.load(std::memory_order::relaxed); + if (oldValue != UnlockedValue) { + return false; + } + return TryAcquireWriter(); +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace NYT::NThreading + diff --git a/library/cpp/yt/threading/writer_starving_rw_spin_lock.cpp b/library/cpp/yt/threading/writer_starving_rw_spin_lock.cpp new file mode 100644 index 00000000000..74c9f59db13 --- /dev/null +++ b/library/cpp/yt/threading/writer_starving_rw_spin_lock.cpp @@ -0,0 +1,25 @@ +#include "writer_starving_rw_spin_lock.h" + +namespace NYT::NThreading { + +//////////////////////////////////////////////////////////////////////////////// + +void TWriterStarvingRWSpinLock::AcquireReaderSlow() noexcept +{ + TSpinWait spinWait(Location_, ESpinLockActivityKind::Read); + while (!TryAndTryAcquireReader()) { + spinWait.Wait(); + } +} + +void TWriterStarvingRWSpinLock::AcquireWriterSlow() noexcept +{ + TSpinWait spinWait(Location_, ESpinLockActivityKind::Write); + while (!TryAndTryAcquireWriter()) { + spinWait.Wait(); + } +} + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace NYT::NThreading diff --git a/library/cpp/yt/threading/writer_starving_rw_spin_lock.h b/library/cpp/yt/threading/writer_starving_rw_spin_lock.h new file mode 100644 index 00000000000..8a456afe21b --- /dev/null +++ b/library/cpp/yt/threading/writer_starving_rw_spin_lock.h @@ -0,0 +1,115 @@ +#pragma once + +#include "public.h" +#include "rw_spin_lock.h" +#include "spin_lock_base.h" +#include "spin_lock_count.h" + +#include <library/cpp/yt/memory/public.h> + +#include <util/system/rwlock.h> + +#include <atomic> + +namespace NYT::NThreading { + +//////////////////////////////////////////////////////////////////////////////// + +// TODO(pavook): deprecate it. + +//! Single-writer multiple-readers spin lock. +/*! + * Reader-side calls are pretty cheap. + * WARNING: The lock is unfair, and readers can starve writers. See rw_spin_lock.h for a writer-prioritized lock. + * WARNING: Never use the bare lock if forks are possible: see fork_aware_rw_spin_lock.h for a fork-safe lock. + * Unlike rw_spin_lock.h, reader-side is reentrant here: it is possible to acquire the **reader** lock multiple times + * even in the single thread. + * This doesn't mean you should do it: in fact, you shouldn't: use separate locks for separate entities. + * If you see this class in your code, try migrating to the proper rw_spin_lock.h after ensuring you don't rely on + * reentrant locking. + */ +class TWriterStarvingRWSpinLock + : public TSpinLockBase +{ +public: + using TSpinLockBase::TSpinLockBase; + + //! Acquires the reader lock. + /*! + * Optimized for the case of read-intensive workloads. + * Cheap (just one atomic increment and no spinning if no writers are present). + * Don't use this call if forks are possible: forking at some + * intermediate point inside #AcquireReader may corrupt the lock state and + * leave lock forever stuck for the child process. + */ + void AcquireReader() noexcept; + //! Tries acquiring the reader lock; see #AcquireReader. + //! Returns |true| on success. + bool TryAcquireReader() noexcept; + //! Releases the reader lock. + /*! + * Cheap (just one atomic decrement). + */ + void ReleaseReader() noexcept; + + //! Acquires the writer lock. + /*! + * Rather cheap (just one CAS). + */ + void AcquireWriter() noexcept; + //! Tries acquiring the writer lock; see #AcquireWriter. + //! Returns |true| on success. + bool TryAcquireWriter() noexcept; + //! Releases the writer lock. + /*! + * Cheap (just one atomic store). + */ + void ReleaseWriter() noexcept; + + //! Returns true if the lock is taken (either by a reader or writer). + /*! + * This is inherently racy. + * Only use for debugging and diagnostic purposes. + */ + bool IsLocked() const noexcept; + + //! Returns true if the lock is taken by reader. + /*! + * This is inherently racy. + * Only use for debugging and diagnostic purposes. + */ + bool IsLockedByReader() const noexcept; + + //! Returns true if the lock is taken by writer. + /*! + * This is inherently racy. + * Only use for debugging and diagnostic purposes. + */ + bool IsLockedByWriter() const noexcept; + +private: + using TValue = ui32; + static constexpr TValue UnlockedValue = 0; + static constexpr TValue WriterMask = 1; + static constexpr TValue ReaderDelta = 2; + + std::atomic<TValue> Value_ = UnlockedValue; + + + bool TryAndTryAcquireReader() noexcept; + bool TryAndTryAcquireWriter() noexcept; + + void AcquireReaderSlow() noexcept; + void AcquireWriterSlow() noexcept; +}; + +REGISTER_TRACKED_SPIN_LOCK_CLASS(TWriterStarvingRWSpinLock) + +//////////////////////////////////////////////////////////////////////////////// + +} // namespace NYT::NThreading + +#define WRITER_STARVING_RW_SPIN_LOCK_INL_H_ +#include "writer_starving_rw_spin_lock-inl.h" +#undef WRITER_STARVING_RW_SPIN_LOCK_INL_H_ + diff --git a/library/cpp/yt/threading/ya.make b/library/cpp/yt/threading/ya.make index cc11e7974ef..d25f0a70681 100644 --- a/library/cpp/yt/threading/ya.make +++ b/library/cpp/yt/threading/ya.make @@ -18,6 +18,7 @@ SRCS( spin_lock.cpp spin_wait.cpp spin_wait_hook.cpp + writer_starving_rw_spin_lock.cpp ) PEERDIR( diff --git a/yql/essentials/core/common_opt/yql_co_flow2.cpp b/yql/essentials/core/common_opt/yql_co_flow2.cpp index 7d521b2ac47..8ebb5243ffa 100644 --- a/yql/essentials/core/common_opt/yql_co_flow2.cpp +++ b/yql/essentials/core/common_opt/yql_co_flow2.cpp @@ -2027,7 +2027,100 @@ void RegisterCoFlowCallables2(TCallableOptimizerMap& map) { return ret; } - return node; + // Add PruneKeys to EquiJoin + static const char optName[] = "EmitPruneKeys"; + if (!IsOptimizerEnabled<optName>(*optCtx.Types) || IsOptimizerDisabled<optName>(*optCtx.Types)) { + return node; + } + auto equiJoin = TCoEquiJoin(node); + if (HasSetting(equiJoin.Arg(equiJoin.ArgCount() - 1).Ref(), "prune_keys_added")) { + return node; + } + + THashMap<TStringBuf, THashSet<TStringBuf>> columnsForPruneKeysExtractor; + GetPruneKeysColumnsForJoinLeaves(equiJoin.Arg(equiJoin.ArgCount() - 2).Cast<TCoEquiJoinTuple>(), columnsForPruneKeysExtractor); + + TExprNode::TListType children; + bool hasChanges = false; + for (size_t i = 0; i + 2 < equiJoin.ArgCount(); ++i) { + auto child = equiJoin.Arg(i).Cast<TCoEquiJoinInput>(); + auto list = child.List(); + auto scope = child.Scope(); + + if (!scope.Ref().IsAtom()) { + children.push_back(equiJoin.Arg(i).Ptr()); + continue; + } + + auto itemNames = columnsForPruneKeysExtractor.find(scope.Ref().Content()); + if (itemNames == columnsForPruneKeysExtractor.end() || itemNames->second.empty()) { + children.push_back(equiJoin.Arg(i).Ptr()); + continue; + } + + if (auto distinct = list.Ref().GetConstraint<TDistinctConstraintNode>()) { + if (distinct->ContainsCompleteSet(std::vector<std::string_view>(itemNames->second.cbegin(), itemNames->second.cend()))) { + children.push_back(equiJoin.Arg(i).Ptr()); + continue; + } + } + + bool isOrdered = false; + if (auto sorted = list.Ref().GetConstraint<TSortedConstraintNode>()) { + for (const auto& item : sorted->GetContent()) { + size_t foundItemNamesCount = 0; + for (const auto& path : item.first) { + if (itemNames->second.contains(path.front())) { + foundItemNamesCount++; + } + } + if (foundItemNamesCount == itemNames->second.size()) { + isOrdered = true; + break; + } + } + } + + auto pruneKeysCallable = isOrdered ? "PruneAdjacentKeys" : "PruneKeys"; + YQL_CLOG(DEBUG, Core) << "Add " << pruneKeysCallable << " to EquiJoin input #" << i << ", label " << scope.Ref().Content(); + children.push_back(ctx.Builder(child.Pos()) + .List() + .Callable(0, pruneKeysCallable) + .Add(0, list.Ptr()) + .Lambda(1) + .Param("item") + .List(0) + .Do([&](TExprNodeBuilder& parent) -> TExprNodeBuilder & { + ui32 i = 0; + for (const auto& column : itemNames->second) { + parent.Callable(i++, "Member") + .Arg(0, "item") + .Atom(1, column) + .Seal(); + } + return parent; + }) + .Seal() + .Seal() + .Seal() + .Add(1, scope.Ptr()) + .Seal() + .Build()); + hasChanges = true; + } + + if (!hasChanges) { + return node; + } + + children.push_back(equiJoin.Arg(equiJoin.ArgCount() - 2).Ptr()); + children.push_back(AddSetting( + equiJoin.Arg(equiJoin.ArgCount() - 1).Ref(), + equiJoin.Arg(equiJoin.ArgCount() - 1).Pos(), + "prune_keys_added", + nullptr, + ctx)); + return ctx.ChangeChildren(*node, std::move(children)); }; map["ExtractMembers"] = [](const TExprNode::TPtr& node, TExprContext& ctx, TOptimizeContext& optCtx) { diff --git a/yql/essentials/core/peephole_opt/yql_opt_peephole_physical.cpp b/yql/essentials/core/peephole_opt/yql_opt_peephole_physical.cpp index 5755bf1dbf5..28ee1ead8f5 100644 --- a/yql/essentials/core/peephole_opt/yql_opt_peephole_physical.cpp +++ b/yql/essentials/core/peephole_opt/yql_opt_peephole_physical.cpp @@ -2760,19 +2760,11 @@ TExprNode::TPtr ExpandListHas(const TExprNode::TPtr& input, TExprContext& ctx) { return RewriteSearchByKeyForTypesMismatch<true, true>(input, ctx); } -TExprNode::TPtr ExpandPruneAdjacentKeys(const TExprNode::TPtr& input, TExprContext& ctx) { - const auto type = input->Head().GetTypeAnn(); - const auto& keyExtractorLambda = input->ChildRef(1); - - YQL_ENSURE(type->GetKind() == ETypeAnnotationKind::List || type->GetKind() == ETypeAnnotationKind::Stream); - - const auto elemType = type->GetKind() == ETypeAnnotationKind::List - ? type->Cast<TListExprType>()->GetItemType() - : type->Cast<TStreamExprType>()->GetItemType(); - const auto optionalElemType = *ctx.MakeType<TOptionalExprType>(elemType); +TExprNode::TPtr ExpandPruneAdjacentKeys(const TExprNode::TPtr& input, TExprContext& ctx, TTypeAnnotationContext& /*typesCtx*/) { + const auto& keyExtractorLambda = input->ChildPtr(1); YQL_CLOG(DEBUG, CorePeepHole) << "Expand " << input->Content(); - return ctx.Builder(input->Pos()) + return KeepConstraints(ctx.Builder(input->Pos()) .Callable("OrderedFlatMap") .Callable(0, "Fold1Map") .Add(0, input->HeadPtr()) @@ -2799,7 +2791,11 @@ TExprNode::TPtr ExpandPruneAdjacentKeys(const TExprNode::TPtr& input, TExprConte .Seal() .Seal() .Callable(1, "Nothing") - .Add(0, ExpandType(input->Pos(), optionalElemType, ctx)) + .Callable(0, "OptionalType") + .Callable(0, "TypeOf") + .Arg(0, "item") + .Seal() + .Seal() .Seal() .Callable(2, "Just") .Arg(0, "item") @@ -2814,13 +2810,16 @@ TExprNode::TPtr ExpandPruneAdjacentKeys(const TExprNode::TPtr& input, TExprConte .Arg(0, "item") .Seal() .Seal() - .Build(); + .Build(), *input, ctx); } -TExprNode::TPtr ExpandPruneKeys(const TExprNode::TPtr& input, TExprContext& ctx) { +TExprNode::TPtr ExpandPruneKeys(const TExprNode::TPtr& input, TExprContext& ctx, TTypeAnnotationContext& typesCtx) { const auto type = input->Head().GetTypeAnn(); - const auto& keyExtractorLambda = input->ChildRef(1); - YQL_ENSURE(type->GetKind() == ETypeAnnotationKind::List || type->GetKind() == ETypeAnnotationKind::Stream); + auto keyExtractorLambda = input->ChildPtr(1); + + YQL_ENSURE(type->GetKind() == ETypeAnnotationKind::Flow + || type->GetKind() == ETypeAnnotationKind::List + || type->GetKind() == ETypeAnnotationKind::Stream); auto initHandler = ctx.Builder(input->Pos()) .Lambda() @@ -2865,6 +2864,39 @@ TExprNode::TPtr ExpandPruneKeys(const TExprNode::TPtr& input, TExprContext& ctx) .Seal() .Build(); } else { + // Slight copy of GetDictionaryKeyTypes to check if keyExtractorLambda result type is complicated + // mkql CombineCore supports only simple types; for others we should add pickling + bool keyExtractorLambdaShouldBePickled = false; + auto itemType = keyExtractorLambda->GetTypeAnn(); + if (itemType->GetKind() == ETypeAnnotationKind::Optional) { + itemType = itemType->Cast<TOptionalExprType>()->GetItemType(); + } + + if (itemType->GetKind() == ETypeAnnotationKind::Tuple) { + auto tuple = itemType->Cast<TTupleExprType>(); + for (const auto& item : tuple->GetItems()) { + if (!IsDataOrOptionalOfData(item)) { + keyExtractorLambdaShouldBePickled = true; + break; + } + } + } else if (itemType->GetKind() != ETypeAnnotationKind::Data) { + keyExtractorLambdaShouldBePickled = true; + } + + if (keyExtractorLambdaShouldBePickled) { + keyExtractorLambda = ctx.Builder(input->Pos()) + .Lambda() + .Param("item") + .Callable(0, "StablePickle") + .Apply(0, keyExtractorLambda) + .With(0, "item") + .Seal() + .Seal() + .Seal() + .Build(); + } + return ctx.Builder(input->Pos()) .Callable("CombineCore") .Add(0, input->HeadPtr()) @@ -2872,6 +2904,7 @@ TExprNode::TPtr ExpandPruneKeys(const TExprNode::TPtr& input, TExprContext& ctx) .Add(2, initHandler) .Add(3, updateHandler) .Add(4, finishHandler) + .Atom(5, ToString(typesCtx.PruneKeysMemLimit)) .Seal() .Build(); } @@ -8906,8 +8939,6 @@ struct TPeepHoleRules { {"CheckedDiv", &ExpandCheckedDiv}, {"CheckedMod", &ExpandCheckedMod}, {"CheckedMinus", &ExpandCheckedMinus}, - {"PruneAdjacentKeys", &ExpandPruneAdjacentKeys}, - {"PruneKeys", &ExpandPruneKeys}, {"JsonValue", &ExpandJsonValue}, {"JsonExists", &ExpandJsonExists}, {"EmptyIterator", &DropDependsOnFromEmptyIterator}, @@ -8926,6 +8957,8 @@ struct TPeepHoleRules { {"CostsOf", &ExpandCostsOf}, {"JsonQuery", &ExpandJsonQuery}, {"MatchRecognize", &ExpandMatchRecognize}, + {"PruneAdjacentKeys", &ExpandPruneAdjacentKeys}, + {"PruneKeys", &ExpandPruneKeys}, {"CalcOverWindow", &ExpandCalcOverWindow}, {"CalcOverSessionWindow", &ExpandCalcOverWindow}, {"CalcOverWindowGroup", &ExpandCalcOverWindow}, diff --git a/yql/essentials/core/yql_expr_constraint.cpp b/yql/essentials/core/yql_expr_constraint.cpp index e67557d9645..7fcf560a2cf 100644 --- a/yql/essentials/core/yql_expr_constraint.cpp +++ b/yql/essentials/core/yql_expr_constraint.cpp @@ -822,9 +822,23 @@ private: } if constexpr (Adjacent) { - return CopyAllFrom<0>(input, output, ctx); + if (const auto status = CopyAllFrom<0>(input, output, ctx); status != TStatus::Ok) { + return status; + } + + TPartOfConstraintBase::TSetType keys = GetPathsToKeys<true>(input->Child(1)->Tail(), input->Child(1)->Head().Head()); + TPartOfConstraintBase::TSetOfSetsType uniqueKeys; + for (const auto& elem : keys) { + uniqueKeys.insert(TPartOfConstraintBase::TSetType{elem}); + } + if (!keys.empty()) { + input->AddConstraint(ctx.MakeConstraint<TUniqueConstraintNode>(TUniqueConstraintNode::TContentType{uniqueKeys})); + input->AddConstraint(ctx.MakeConstraint<TDistinctConstraintNode>(TDistinctConstraintNode::TContentType{uniqueKeys})); + } + return TStatus::Ok; } - return FromFirst<TEmptyConstraintNode, TUniqueConstraintNode, TPartOfUniqueConstraintNode, TDistinctConstraintNode, TPartOfDistinctConstraintNode>(input, output, ctx); + + return FromFirst<TEmptyConstraintNode>(input, output, ctx); } template<class TConstraint> diff --git a/yql/essentials/core/yql_join.cpp b/yql/essentials/core/yql_join.cpp index 7cca45604d1..5ce3fba268b 100644 --- a/yql/essentials/core/yql_join.cpp +++ b/yql/essentials/core/yql_join.cpp @@ -820,6 +820,10 @@ IGraphTransformer::TStatus ValidateEquiJoinOptions(TPositionHandle positionHandl // do nothing } else if (optionName == "multiple_joins") { // do nothing + } else if (optionName == "prune_keys_added") { + if (!EnsureTupleSize(*child, 1, ctx)) { + return IGraphTransformer::TStatus::Error; + } } else { YQL_ENSURE(false, "Cached join option '" << optionName << "' not handled"); } @@ -2007,7 +2011,7 @@ void GatherJoinInputs(const TExprNode::TPtr& expr, const TExprNode& row, } bool IsCachedJoinOption(TStringBuf name) { - static THashSet<TStringBuf> CachedJoinOptions = {"preferred_sort", "cbo_passed", "multiple_joins"}; + static THashSet<TStringBuf> CachedJoinOptions = {"preferred_sort", "cbo_passed", "multiple_joins", "prune_keys_added"}; return CachedJoinOptions.contains(name); } @@ -2016,4 +2020,31 @@ bool IsCachedJoinLinkOption(TStringBuf name) { return CachedJoinLinkOptions.contains(name); } +void GetPruneKeysColumnsForJoinLeaves(const TCoEquiJoinTuple& joinTree, THashMap<TStringBuf, THashSet<TStringBuf>>& columnsForPruneKeysExtractor) { + auto settings = GetEquiJoinLinkSettings(joinTree.Options().Ref()); + TStringBuf joinKind = joinTree.Type().Value(); + + auto left = joinTree.LeftScope(); + if (!left.Maybe<TCoAtom>()) { + GetPruneKeysColumnsForJoinLeaves(left.Cast<TCoEquiJoinTuple>(), columnsForPruneKeysExtractor); + } else { + if (joinKind == "RightSemi" || joinKind == "RightOnly" || settings.LeftHints.contains("any")) { + if (!settings.LeftHints.contains("unique")) { + CollectEquiJoinKeyColumnsFromLeaf(joinTree.LeftKeys().Ref(), columnsForPruneKeysExtractor); + } + } + } + + auto right = joinTree.RightScope(); + if (!right.Maybe<TCoAtom>()) { + GetPruneKeysColumnsForJoinLeaves(right.Cast<TCoEquiJoinTuple>(), columnsForPruneKeysExtractor); + } else { + if (joinKind == "LeftSemi" || joinKind == "LeftOnly" || settings.RightHints.contains("any")) { + if (!settings.RightHints.contains("unique")) { + CollectEquiJoinKeyColumnsFromLeaf(joinTree.RightKeys().Ref(), columnsForPruneKeysExtractor); + } + } + } +} + } // namespace NYql diff --git a/yql/essentials/core/yql_join.h b/yql/essentials/core/yql_join.h index a313aceefa3..26b417cf542 100644 --- a/yql/essentials/core/yql_join.h +++ b/yql/essentials/core/yql_join.h @@ -180,4 +180,6 @@ void GatherJoinInputs(const TExprNode::TPtr& expr, const TExprNode& row, bool IsCachedJoinOption(TStringBuf name); bool IsCachedJoinLinkOption(TStringBuf name); +void GetPruneKeysColumnsForJoinLeaves(const NNodes::TCoEquiJoinTuple& joinTree, THashMap<TStringBuf, THashSet<TStringBuf>>& columnsForPruneKeysExtractor); + } diff --git a/yql/essentials/core/yql_type_annotation.h b/yql/essentials/core/yql_type_annotation.h index bc09c963f4b..a43650a01f4 100644 --- a/yql/essentials/core/yql_type_annotation.h +++ b/yql/essentials/core/yql_type_annotation.h @@ -452,6 +452,7 @@ struct TTypeAnnotationContext: public TThrRefBase { THashSet<TString> PeepholeFlags; bool StreamLookupJoin = false; ui32 MaxAggPushdownPredicates = 6; // algorithm complexity is O(2^N) + ui32 PruneKeysMemLimit = 128 * 1024 * 1024; TMaybe<TColumnOrder> LookupColumnOrder(const TExprNode& node) const; IGraphTransformer::TStatus SetColumnOrder(const TExprNode& node, const TColumnOrder& columnOrder, TExprContext& ctx); diff --git a/yql/essentials/data/language/pragmas_opensource.json b/yql/essentials/data/language/pragmas_opensource.json index 26cc1c777a0..a4a3ad73364 100644 --- a/yql/essentials/data/language/pragmas_opensource.json +++ b/yql/essentials/data/language/pragmas_opensource.json @@ -1 +1 @@ -[{"name":"yt.Annotations"},{"name":"yt.ApplyStoredConstraints"},{"name":"yt.Auth"},{"name":"yt.AutoMerge"},{"name":"yt.BatchListFolderConcurrency"},{"name":"yt.BinaryExpirationInterval"},{"name":"yt.BinaryTmpFolder"},{"name":"yt.BlockMapJoin"},{"name":"yt.BlockReaderSupportedDataTypes"},{"name":"yt.BlockReaderSupportedTypes"},{"name":"yt.BufferRowCount"},{"name":"yt.ClientMapTimeout"},{"name":"yt.ColumnGroupMode"},{"name":"yt.CombineCoreLimit"},{"name":"yt.CommonJoinCoreLimit"},{"name":"yt.CompactForDistinct"},{"name":"yt.CoreDumpPath"},{"name":"yt.DQRPCReaderInflight"},{"name":"yt.DQRPCReaderTimeout"},{"name":"yt.DataSizePerJob"},{"name":"yt.DataSizePerMapJob"},{"name":"yt.DataSizePerPartition"},{"name":"yt.DataSizePerSortJob"},{"name":"yt.DefaultCalcMemoryLimit"},{"name":"yt.DefaultCluster"},{"name":"yt.DefaultLocalityTimeout"},{"name":"yt.DefaultMapSelectivityFactor"},{"name":"yt.DefaultMaxJobFails"},{"name":"yt.DefaultMemoryDigestLowerBound"},{"name":"yt.DefaultMemoryLimit"},{"name":"yt.DefaultMemoryReserveFactor"},{"name":"yt.DefaultOperationWeight"},{"name":"yt.DefaultRuntimeCluster"},{"name":"yt.Description"},{"name":"yt.DisableFuseOperations"},{"name":"yt.DisableJobSplitting"},{"name":"yt.DisableOptimizers"},{"name":"yt.DockerImage"},{"name":"yt.DqPruneKeyFilterLambda"},{"name":"yt.DropUnusedKeysFromKeyFilter"},{"name":"yt.EnableDynamicStoreReadInDQ"},{"name":"yt.EnableFuseMapToMapReduce"},{"name":"yt.EnforceJobUtc"},{"name":"yt.ErasureCodecCpu"},{"name":"yt.ErasureCodecCpuForDq"},{"name":"yt.EvaluationTableSizeLimit"},{"name":"yt.ExpirationDeadline"},{"name":"yt.ExpirationInterval"},{"name":"yt.ExtendTableLimit"},{"name":"yt.ExtendedStatsMaxChunkCount"},{"name":"yt.ExternalTx"},{"name":"yt.ExtraTmpfsSize"},{"name":"yt.FileCacheTtl"},{"name":"yt.FmrOperationSpec"},{"name":"yt.FolderInlineDataLimit"},{"name":"yt.FolderInlineItemsLimit"},{"name":"yt.ForceInferSchema"},{"name":"yt.ForceJobSizeAdjuster"},{"name":"yt.ForceTmpSecurity"},{"name":"yt.GeobaseDownloadUrl"},{"name":"yt.HybridDqDataSizeLimitForOrdered"},{"name":"yt.HybridDqDataSizeLimitForUnordered"},{"name":"yt.HybridDqExecution"},{"name":"yt.HybridDqExecutionFallback"},{"name":"yt.IgnoreTypeV3"},{"name":"yt.IgnoreWeakSchema"},{"name":"yt.IgnoreYamrDsv"},{"name":"yt.InferSchema"},{"name":"yt.InferSchemaMode"},{"name":"yt.InferSchemaTableCountThreshold"},{"name":"yt.InflightTempTablesLimit"},{"name":"yt.IntermediateAccount"},{"name":"yt.IntermediateDataMedium"},{"name":"yt.IntermediateReplicationFactor"},{"name":"yt.JavascriptCpu"},{"name":"yt.JobBlockInput"},{"name":"yt.JobBlockInputSupportedDataTypes"},{"name":"yt.JobBlockInputSupportedTypes"},{"name":"yt.JobBlockOutput"},{"name":"yt.JobBlockOutputSupportedDataTypes"},{"name":"yt.JobBlockOutputSupportedTypes"},{"name":"yt.JobBlockTableContent"},{"name":"yt.JobEnv"},{"name":"yt.JoinAllowColumnRenames"},{"name":"yt.JoinCollectColumnarStatistics"},{"name":"yt.JoinColumnarStatisticsFetcherMode"},{"name":"yt.JoinCommonUseMapMultiOut"},{"name":"yt.JoinEnableStarJoin"},{"name":"yt.JoinMergeForce"},{"name":"yt.JoinMergeReduceJobMaxSize"},{"name":"yt.JoinMergeSetTopLevelFullSort"},{"name":"yt.JoinMergeTablesLimit"},{"name":"yt.JoinMergeUnsortedFactor"},{"name":"yt.JoinMergeUseSmallAsPrimary"},{"name":"yt.JoinUseColumnarStatistics"},{"name":"yt.JoinWaitAllInputs"},{"name":"yt.KeepTempTables"},{"name":"yt.KeyFilterForStartsWith"},{"name":"yt.LLVMMemSize"},{"name":"yt.LLVMNodeCountLimit"},{"name":"yt.LLVMPerNodeMemSize"},{"name":"yt.LayerPaths"},{"name":"yt.LocalCalcLimit"},{"name":"yt.LookupJoinLimit"},{"name":"yt.LookupJoinMaxRows"},{"name":"yt.MapJoinLimit"},{"name":"yt.MapJoinShardCount"},{"name":"yt.MapJoinShardMinRows"},{"name":"yt.MapJoinUseFlow"},{"name":"yt.MapLocalityTimeout"},{"name":"yt.MaxChunksForDqRead"},{"name":"yt.MaxColumnGroups"},{"name":"yt.MaxCpuUsageToFuseMultiOuts"},{"name":"yt.MaxExtraJobMemoryToFuseOperations"},{"name":"yt.MaxInputTables"},{"name":"yt.MaxInputTablesForSortedMerge"},{"name":"yt.MaxJobCount"},{"name":"yt.MaxKeyRangeCount"},{"name":"yt.MaxKeyWeight"},{"name":"yt.MaxOperationFiles"},{"name":"yt.MaxOutputTables"},{"name":"yt.MaxReplicationFactorToFuseMultiOuts"},{"name":"yt.MaxReplicationFactorToFuseOperations"},{"name":"yt.MaxRowWeight"},{"name":"yt.MaxSpeculativeJobCountPerTask"},{"name":"yt.MergeAdjacentPointRanges"},{"name":"yt.MinColumnGroupSize"},{"name":"yt.MinLocalityInputDataWeight"},{"name":"yt.MinPublishedAvgChunkSize"},{"name":"yt.MinTempAvgChunkSize"},{"name":"yt.NativeYtTypeCompatibility"},{"name":"yt.NetworkProject"},{"name":"yt.NightlyCompress"},{"name":"yt.OperationReaders"},{"name":"yt.OperationSpec"},{"name":"yt.OptimizeFor"},{"name":"yt.Owners"},{"name":"yt.ParallelOperationsLimit"},{"name":"yt.PartitionByConstantKeysViaMap"},{"name":"yt.Pool"},{"name":"yt.PoolTrees"},{"name":"yt.PrimaryMedium"},{"name":"yt.PruneKeyFilterLambda"},{"name":"yt.PruneQLFilterLambda"},{"name":"yt.PublishedAutoMerge"},{"name":"yt.PublishedCompressionCodec"},{"name":"yt.PublishedErasureCodec"},{"name":"yt.PublishedMedia"},{"name":"yt.PublishedPrimaryMedium"},{"name":"yt.PublishedReplicationFactor"},{"name":"yt.PythonCpu"},{"name":"yt.QueryCacheChunkLimit"},{"name":"yt.QueryCacheIgnoreTableRevision"},{"name":"yt.QueryCacheMode"},{"name":"yt.QueryCacheSalt"},{"name":"yt.QueryCacheTtl"},{"name":"yt.QueryCacheUseExpirationTimeout"},{"name":"yt.QueryCacheUseForCalc"},{"name":"yt.ReduceLocalityTimeout"},{"name":"yt.ReleaseTempData"},{"name":"yt.ReportEquiJoinStats"},{"name":"yt.RuntimeCluster"},{"name":"yt.RuntimeClusterSelection"},{"name":"yt.SamplingIoBlockSize"},{"name":"yt.SchedulingTag"},{"name":"yt.SchedulingTagFilter"},{"name":"yt.ScriptCpu"},{"name":"yt.SortLocalityTimeout"},{"name":"yt.StartedBy"},{"name":"yt.StaticPool"},{"name":"yt.SuspendIfAccountLimitExceeded"},{"name":"yt.SwitchLimit"},{"name":"yt.TableContentColumnarStatistics"},{"name":"yt.TableContentCompressLevel"},{"name":"yt.TableContentDeliveryMode"},{"name":"yt.TableContentLocalExecution"},{"name":"yt.TableContentMaxChunksForNativeDelivery"},{"name":"yt.TableContentMaxInputTables"},{"name":"yt.TableContentMinAvgChunkSize"},{"name":"yt.TableContentTmpFolder"},{"name":"yt.TableContentUseSkiff"},{"name":"yt.TablesTmpFolder"},{"name":"yt.TempTablesTtl"},{"name":"yt.TemporaryAutoMerge"},{"name":"yt.TemporaryCompressionCodec"},{"name":"yt.TemporaryErasureCodec"},{"name":"yt.TemporaryMedia"},{"name":"yt.TemporaryPrimaryMedium"},{"name":"yt.TemporaryReplicationFactor"},{"name":"yt.TentativePoolTrees"},{"name":"yt.TentativeTreeEligibilityMaxJobDurationRatio"},{"name":"yt.TentativeTreeEligibilityMinJobDuration"},{"name":"yt.TentativeTreeEligibilitySampleJobCount"},{"name":"yt.TmpFolder"},{"name":"yt.TopSortMaxLimit"},{"name":"yt.TopSortRowMultiplierPerJob"},{"name":"yt.TopSortSizePerJob"},{"name":"yt.UseAggPhases"},{"name":"yt.UseColumnGroupsFromInputTables"},{"name":"yt.UseColumnarStatistics"},{"name":"yt.UseDefaultTentativePoolTrees"},{"name":"yt.UseFlow"},{"name":"yt.UseIntermediateSchema"},{"name":"yt.UseIntermediateStreams"},{"name":"yt.UseNativeDescSort"},{"name":"yt.UseNativeYtTypes"},{"name":"yt.UseNewPredicateExtraction"},{"name":"yt.UsePartitionsByKeysForFinalAgg"},{"name":"yt.UseQLFilter"},{"name":"yt.UseRPCReaderInDQ"},{"name":"yt.UseSkiff"},{"name":"yt.UseSystemColumns"},{"name":"yt.UseTmpfs"},{"name":"yt.UseTypeV2"},{"name":"yt.UseYqlRowSpecCompactForm"},{"name":"yt.UserSlots"},{"name":"yt.ViewIsolation"},{"name":"yt.WideFlowLimit"},{"name":"dq.AggregateStatsByStage"},{"name":"dq.AnalyticsHopping"},{"name":"dq.AnalyzeQuery"},{"name":"dq.ChannelBufferSize"},{"name":"dq.ChunkSizeLimit"},{"name":"dq.CollectCoreDumps"},{"name":"dq.ComputeActorType"},{"name":"dq.DataSizePerJob"},{"name":"dq.DisableCheckpoints"},{"name":"dq.DisableLLVMForBlockStages"},{"name":"dq.EnableChannelStats"},{"name":"dq.EnableComputeActor"},{"name":"dq.EnableDqReplicate"},{"name":"dq.EnableFullResultWrite"},{"name":"dq.EnableInsert"},{"name":"dq.EnableSpillingInChannels"},{"name":"dq.EnableSpillingNodes"},{"name":"dq.EnableStrip"},{"name":"dq.ExportStats"},{"name":"dq.FallbackPolicy"},{"name":"dq.HashJoinMode"},{"name":"dq.HashShuffleMaxTasks"},{"name":"dq.HashShuffleTasksRatio"},{"name":"dq.MaxDataSizePerJob"},{"name":"dq.MaxDataSizePerQuery"},{"name":"dq.MaxNetworkRetries"},{"name":"dq.MaxRetries"},{"name":"dq.MaxTasksPerOperation"},{"name":"dq.MaxTasksPerStage"},{"name":"dq.MemoryLimit"},{"name":"dq.OptLLVM"},{"name":"dq.OutputChunkMaxSize"},{"name":"dq.ParallelOperationsLimit"},{"name":"dq.PingTimeoutMs"},{"name":"dq.PullRequestTimeoutMs"},{"name":"dq.QueryTimeout"},{"name":"dq.RetryBackoffMs"},{"name":"dq.Scheduler"},{"name":"dq.SpillingEngine"},{"name":"dq.SplitStageOnDqReplicate"},{"name":"dq.TaskRunnerStats"},{"name":"dq.UseAggPhases"},{"name":"dq.UseBlockReader"},{"name":"dq.UseFastPickleTransport"},{"name":"dq.UseFinalizeByKey"},{"name":"dq.UseGraceJoinCoreForMap"},{"name":"dq.UseOOBTransport"},{"name":"dq.UseSimpleYtReader"},{"name":"dq.UseWideBlockChannels"},{"name":"dq.UseWideChannels"},{"name":"dq.WatermarksEnableIdlePartitions"},{"name":"dq.WatermarksGranularityMs"},{"name":"dq.WatermarksLateArrivalDelayMs"},{"name":"dq.WatermarksMode"},{"name":"dq.WorkerFilter"},{"name":"dq.WorkersPerOperation"},{"name":"AllowDotInAlias"},{"name":"AllowUnnamedColumns"},{"name":"AnsiCurrentRow"},{"name":"AnsiImplicitCrossJoin"},{"name":"AnsiInForEmptyOrNullableItemsCollections"},{"name":"AnsiLike"},{"name":"AnsiOptionalAs"},{"name":"AnsiRankForNullableKeys"},{"name":"AutoCommit"},{"name":"BlockEngine"},{"name":"BlockEngineEnable"},{"name":"BlockEngineForce"},{"name":"BogousStarInGroupByOverJoin"},{"name":"CheckedOps"},{"name":"ClassicDivision"},{"name":"CoalesceJoinKeysOnQualifiedAll"},{"name":"CompactGroupBy"},{"name":"CompactNamedExprs"},{"name":"CostBasedOptimizer"},{"name":"DataWatermarks"},{"name":"DirectRead"},{"name":"DisableAnsiCurrentRow"},{"name":"DisableAnsiImplicitCrossJoin"},{"name":"DisableAnsiInForEmptyOrNullableItemsCollections"},{"name":"DisableAnsiLike"},{"name":"DisableAnsiOptionalAs"},{"name":"DisableAnsiRankForNullableKeys"},{"name":"DisableBlockEngineEnable"},{"name":"DisableBlockEngineForce"},{"name":"DisableBogousStarInGroupByOverJoin"},{"name":"DisableCoalesceJoinKeysOnQualifiedAll"},{"name":"DisableCompactGroupBy"},{"name":"DisableCompactNamedExprs"},{"name":"DisableDistinctOverKeys"},{"name":"DisableDistinctOverWindow"},{"name":"DisableDqEngineEnable"},{"name":"DisableDqEngineForce"},{"name":"DisableEmitAggApply"},{"name":"DisableEmitStartsWith"},{"name":"DisableEmitTableSource"},{"name":"DisableEmitUnionMerge"},{"name":"DisableFilterPushdownOverJoinOptionalSide"},{"name":"DisableFlexibleTypes"},{"name":"DisableJsonQueryReturnsJsonDocument"},{"name":"DisableOrderedColumns"},{"name":"DisablePullUpFlatMapOverJoin"},{"name":"DisableRegexUseRe2"},{"name":"DisableRotateJoinTree"},{"name":"DisableSeqMode"},{"name":"DisableSimpleColumns"},{"name":"DisableStrictJoinKeyTypes"},{"name":"DisableUnicodeLiterals"},{"name":"DisableUnorderedResult"},{"name":"DisableUnorderedSubqueries"},{"name":"DisableUseBlocks"},{"name":"DisableValidateUnusedExprs"},{"name":"DisableWarnOnAnsiAliasShadowing"},{"name":"DisableWarnUntypedStringLiterals"},{"name":"DiscoveryMode"},{"name":"DistinctOverKeys"},{"name":"DistinctOverWindow"},{"name":"DqEngine"},{"name":"DqEngineEnable"},{"name":"DqEngineForce"},{"name":"EmitAggApply"},{"name":"EmitStartsWith"},{"name":"EmitTableSource"},{"name":"EmitUnionMerge"},{"name":"EnableSystemColumns"},{"name":"Engine"},{"name":"ErrorMsg"},{"name":"FeatureR010"},{"name":"File"},{"name":"FileOption"},{"name":"FilterPushdownOverJoinOptionalSide"},{"name":"FlexibleTypes"},{"name":"Folder"},{"name":"Greetings"},{"name":"GroupByCubeLimit"},{"name":"GroupByLimit"},{"name":"JsonQueryReturnsJsonDocument"},{"name":"Library"},{"name":"OrderedColumns"},{"name":"OverrideLibrary"},{"name":"Package"},{"name":"PackageVersion"},{"name":"PathPrefix"},{"name":"PositionalUnionAll"},{"name":"PqReadBy"},{"name":"PullUpFlatMapOverJoin"},{"name":"RefSelect"},{"name":"RegexUseRe2"},{"name":"ResultRowsLimit"},{"name":"ResultSizeLimit"},{"name":"RotateJoinTree"},{"name":"RuntimeLogLevel"},{"name":"SampleSelect"},{"name":"SeqMode"},{"name":"SimpleColumns"},{"name":"StrictJoinKeyTypes"},{"name":"Udf"},{"name":"UnicodeLiterals"},{"name":"UnorderedResult"},{"name":"UnorderedSubqueries"},{"name":"UseBlocks"},{"name":"UseTablePrefixForEach"},{"name":"ValidateUnusedExprs"},{"name":"WarnOnAnsiAliasShadowing"},{"name":"WarnUnnamedColumns"},{"name":"WarnUntypedStringLiterals"},{"name":"Warning"},{"name":"WarningMsg"},{"name":"yson.AutoConvert"},{"name":"yson.CastToString"},{"name":"yson.DisableCastToString"},{"name":"yson.DisableStrict"},{"name":"yson.Strict"}] +[{"name":"yt.Annotations"},{"name":"yt.ApplyStoredConstraints"},{"name":"yt.Auth"},{"name":"yt.AutoMerge"},{"name":"yt.BatchListFolderConcurrency"},{"name":"yt.BinaryExpirationInterval"},{"name":"yt.BinaryTmpFolder"},{"name":"yt.BlockMapJoin"},{"name":"yt.BlockReaderSupportedDataTypes"},{"name":"yt.BlockReaderSupportedTypes"},{"name":"yt.BufferRowCount"},{"name":"yt.ClientMapTimeout"},{"name":"yt.ColumnGroupMode"},{"name":"yt.CombineCoreLimit"},{"name":"yt.CommonJoinCoreLimit"},{"name":"yt.CompactForDistinct"},{"name":"yt.CoreDumpPath"},{"name":"yt.DQRPCReaderInflight"},{"name":"yt.DQRPCReaderTimeout"},{"name":"yt.DataSizePerJob"},{"name":"yt.DataSizePerMapJob"},{"name":"yt.DataSizePerPartition"},{"name":"yt.DataSizePerSortJob"},{"name":"yt.DefaultCalcMemoryLimit"},{"name":"yt.DefaultCluster"},{"name":"yt.DefaultLocalityTimeout"},{"name":"yt.DefaultMapSelectivityFactor"},{"name":"yt.DefaultMaxJobFails"},{"name":"yt.DefaultMemoryDigestLowerBound"},{"name":"yt.DefaultMemoryLimit"},{"name":"yt.DefaultMemoryReserveFactor"},{"name":"yt.DefaultOperationWeight"},{"name":"yt.DefaultRuntimeCluster"},{"name":"yt.Description"},{"name":"yt.DisableFuseOperations"},{"name":"yt.DisableJobSplitting"},{"name":"yt.DisableOptimizers"},{"name":"yt.DockerImage"},{"name":"yt.DqPruneKeyFilterLambda"},{"name":"yt.DropUnusedKeysFromKeyFilter"},{"name":"yt.EnableDynamicStoreReadInDQ"},{"name":"yt.EnableFuseMapToMapReduce"},{"name":"yt.EnforceJobUtc"},{"name":"yt.ErasureCodecCpu"},{"name":"yt.ErasureCodecCpuForDq"},{"name":"yt.EvaluationTableSizeLimit"},{"name":"yt.ExpirationDeadline"},{"name":"yt.ExpirationInterval"},{"name":"yt.ExtendTableLimit"},{"name":"yt.ExtendedStatsMaxChunkCount"},{"name":"yt.ExternalTx"},{"name":"yt.ExtraTmpfsSize"},{"name":"yt.FileCacheTtl"},{"name":"yt.FmrOperationSpec"},{"name":"yt.FolderInlineDataLimit"},{"name":"yt.FolderInlineItemsLimit"},{"name":"yt.ForceInferSchema"},{"name":"yt.ForceJobSizeAdjuster"},{"name":"yt.ForceTmpSecurity"},{"name":"yt.GeobaseDownloadUrl"},{"name":"yt.HybridDqDataSizeLimitForOrdered"},{"name":"yt.HybridDqDataSizeLimitForUnordered"},{"name":"yt.HybridDqExecution"},{"name":"yt.HybridDqExecutionFallback"},{"name":"yt.IgnoreTypeV3"},{"name":"yt.IgnoreWeakSchema"},{"name":"yt.IgnoreYamrDsv"},{"name":"yt.InferSchema"},{"name":"yt.InferSchemaMode"},{"name":"yt.InferSchemaTableCountThreshold"},{"name":"yt.InflightTempTablesLimit"},{"name":"yt.IntermediateAccount"},{"name":"yt.IntermediateDataMedium"},{"name":"yt.IntermediateReplicationFactor"},{"name":"yt.JavascriptCpu"},{"name":"yt.JobBlockInput"},{"name":"yt.JobBlockInputSupportedDataTypes"},{"name":"yt.JobBlockInputSupportedTypes"},{"name":"yt.JobBlockOutput"},{"name":"yt.JobBlockOutputSupportedDataTypes"},{"name":"yt.JobBlockOutputSupportedTypes"},{"name":"yt.JobBlockTableContent"},{"name":"yt.JobEnv"},{"name":"yt.JoinAllowColumnRenames"},{"name":"yt.JoinCollectColumnarStatistics"},{"name":"yt.JoinColumnarStatisticsFetcherMode"},{"name":"yt.JoinCommonUseMapMultiOut"},{"name":"yt.JoinEnableStarJoin"},{"name":"yt.JoinMergeForce"},{"name":"yt.JoinMergeReduceJobMaxSize"},{"name":"yt.JoinMergeSetTopLevelFullSort"},{"name":"yt.JoinMergeTablesLimit"},{"name":"yt.JoinMergeUnsortedFactor"},{"name":"yt.JoinMergeUseSmallAsPrimary"},{"name":"yt.JoinUseColumnarStatistics"},{"name":"yt.JoinWaitAllInputs"},{"name":"yt.KeepTempTables"},{"name":"yt.KeyFilterForStartsWith"},{"name":"yt.LLVMMemSize"},{"name":"yt.LLVMNodeCountLimit"},{"name":"yt.LLVMPerNodeMemSize"},{"name":"yt.LayerPaths"},{"name":"yt.LocalCalcLimit"},{"name":"yt.LookupJoinLimit"},{"name":"yt.LookupJoinMaxRows"},{"name":"yt.MapJoinLimit"},{"name":"yt.MapJoinShardCount"},{"name":"yt.MapJoinShardMinRows"},{"name":"yt.MapJoinUseFlow"},{"name":"yt.MapLocalityTimeout"},{"name":"yt.MaxChunksForDqRead"},{"name":"yt.MaxColumnGroups"},{"name":"yt.MaxCpuUsageToFuseMultiOuts"},{"name":"yt.MaxExtraJobMemoryToFuseOperations"},{"name":"yt.MaxInputTables"},{"name":"yt.MaxInputTablesForSortedMerge"},{"name":"yt.MaxJobCount"},{"name":"yt.MaxKeyRangeCount"},{"name":"yt.MaxKeyWeight"},{"name":"yt.MaxOperationFiles"},{"name":"yt.MaxOutputTables"},{"name":"yt.MaxReplicationFactorToFuseMultiOuts"},{"name":"yt.MaxReplicationFactorToFuseOperations"},{"name":"yt.MaxRowWeight"},{"name":"yt.MaxSpeculativeJobCountPerTask"},{"name":"yt.MergeAdjacentPointRanges"},{"name":"yt.MinColumnGroupSize"},{"name":"yt.MinLocalityInputDataWeight"},{"name":"yt.MinPublishedAvgChunkSize"},{"name":"yt.MinTempAvgChunkSize"},{"name":"yt.NativeYtTypeCompatibility"},{"name":"yt.NetworkProject"},{"name":"yt.NightlyCompress"},{"name":"yt.OperationReaders"},{"name":"yt.OperationSpec"},{"name":"yt.OptimizeFor"},{"name":"yt.Owners"},{"name":"yt.ParallelOperationsLimit"},{"name":"yt.PartitionByConstantKeysViaMap"},{"name":"yt.Pool"},{"name":"yt.PoolTrees"},{"name":"yt.PrimaryMedium"},{"name":"yt.PruneKeyFilterLambda"},{"name":"yt.PruneQLFilterLambda"},{"name":"yt.PublishedAutoMerge"},{"name":"yt.PublishedCompressionCodec"},{"name":"yt.PublishedErasureCodec"},{"name":"yt.PublishedMedia"},{"name":"yt.PublishedPrimaryMedium"},{"name":"yt.PublishedReplicationFactor"},{"name":"yt.PythonCpu"},{"name":"yt.QueryCacheChunkLimit"},{"name":"yt.QueryCacheIgnoreTableRevision"},{"name":"yt.QueryCacheMode"},{"name":"yt.QueryCacheSalt"},{"name":"yt.QueryCacheTtl"},{"name":"yt.QueryCacheUseExpirationTimeout"},{"name":"yt.QueryCacheUseForCalc"},{"name":"yt.ReduceLocalityTimeout"},{"name":"yt.ReleaseTempData"},{"name":"yt.ReportEquiJoinStats"},{"name":"yt.RuntimeCluster"},{"name":"yt.RuntimeClusterSelection"},{"name":"yt.SamplingIoBlockSize"},{"name":"yt.SchedulingTag"},{"name":"yt.SchedulingTagFilter"},{"name":"yt.ScriptCpu"},{"name":"yt.SortLocalityTimeout"},{"name":"yt.StartedBy"},{"name":"yt.StaticPool"},{"name":"yt.SuspendIfAccountLimitExceeded"},{"name":"yt.SwitchLimit"},{"name":"yt.TableContentColumnarStatistics"},{"name":"yt.TableContentCompressLevel"},{"name":"yt.TableContentDeliveryMode"},{"name":"yt.TableContentLocalExecution"},{"name":"yt.TableContentMaxChunksForNativeDelivery"},{"name":"yt.TableContentMaxInputTables"},{"name":"yt.TableContentMinAvgChunkSize"},{"name":"yt.TableContentTmpFolder"},{"name":"yt.TableContentUseSkiff"},{"name":"yt.TablesTmpFolder"},{"name":"yt.TempTablesTtl"},{"name":"yt.TemporaryAutoMerge"},{"name":"yt.TemporaryCompressionCodec"},{"name":"yt.TemporaryErasureCodec"},{"name":"yt.TemporaryMedia"},{"name":"yt.TemporaryPrimaryMedium"},{"name":"yt.TemporaryReplicationFactor"},{"name":"yt.TentativePoolTrees"},{"name":"yt.TentativeTreeEligibilityMaxJobDurationRatio"},{"name":"yt.TentativeTreeEligibilityMinJobDuration"},{"name":"yt.TentativeTreeEligibilitySampleJobCount"},{"name":"yt.TmpFolder"},{"name":"yt.TopSortMaxLimit"},{"name":"yt.TopSortRowMultiplierPerJob"},{"name":"yt.TopSortSizePerJob"},{"name":"yt.UseAggPhases"},{"name":"yt.UseColumnGroupsFromInputTables"},{"name":"yt.UseColumnarStatistics"},{"name":"yt.UseDefaultTentativePoolTrees"},{"name":"yt.UseFlow"},{"name":"yt.UseIntermediateSchema"},{"name":"yt.UseIntermediateStreams"},{"name":"yt.UseNativeDescSort"},{"name":"yt.UseNativeDynamicTableRead"},{"name":"yt.UseNativeYtTypes"},{"name":"yt.UseNewPredicateExtraction"},{"name":"yt.UsePartitionsByKeysForFinalAgg"},{"name":"yt.UseQLFilter"},{"name":"yt.UseRPCReaderInDQ"},{"name":"yt.UseSkiff"},{"name":"yt.UseSystemColumns"},{"name":"yt.UseTmpfs"},{"name":"yt.UseTypeV2"},{"name":"yt.UseYqlRowSpecCompactForm"},{"name":"yt.UserSlots"},{"name":"yt.ViewIsolation"},{"name":"yt.WideFlowLimit"},{"name":"dq.AggregateStatsByStage"},{"name":"dq.AnalyticsHopping"},{"name":"dq.AnalyzeQuery"},{"name":"dq.ChannelBufferSize"},{"name":"dq.ChunkSizeLimit"},{"name":"dq.CollectCoreDumps"},{"name":"dq.ComputeActorType"},{"name":"dq.DataSizePerJob"},{"name":"dq.DisableCheckpoints"},{"name":"dq.DisableLLVMForBlockStages"},{"name":"dq.EnableChannelStats"},{"name":"dq.EnableComputeActor"},{"name":"dq.EnableDqReplicate"},{"name":"dq.EnableFullResultWrite"},{"name":"dq.EnableInsert"},{"name":"dq.EnableSpillingInChannels"},{"name":"dq.EnableSpillingNodes"},{"name":"dq.EnableStrip"},{"name":"dq.ExportStats"},{"name":"dq.FallbackPolicy"},{"name":"dq.HashJoinMode"},{"name":"dq.HashShuffleMaxTasks"},{"name":"dq.HashShuffleTasksRatio"},{"name":"dq.MaxDataSizePerJob"},{"name":"dq.MaxDataSizePerQuery"},{"name":"dq.MaxNetworkRetries"},{"name":"dq.MaxRetries"},{"name":"dq.MaxTasksPerOperation"},{"name":"dq.MaxTasksPerStage"},{"name":"dq.MemoryLimit"},{"name":"dq.OptLLVM"},{"name":"dq.OutputChunkMaxSize"},{"name":"dq.ParallelOperationsLimit"},{"name":"dq.PingTimeoutMs"},{"name":"dq.PullRequestTimeoutMs"},{"name":"dq.QueryTimeout"},{"name":"dq.RetryBackoffMs"},{"name":"dq.Scheduler"},{"name":"dq.SpillingEngine"},{"name":"dq.SplitStageOnDqReplicate"},{"name":"dq.TaskRunnerStats"},{"name":"dq.UseAggPhases"},{"name":"dq.UseBlockReader"},{"name":"dq.UseFastPickleTransport"},{"name":"dq.UseFinalizeByKey"},{"name":"dq.UseGraceJoinCoreForMap"},{"name":"dq.UseOOBTransport"},{"name":"dq.UseSimpleYtReader"},{"name":"dq.UseWideBlockChannels"},{"name":"dq.UseWideChannels"},{"name":"dq.WatermarksEnableIdlePartitions"},{"name":"dq.WatermarksGranularityMs"},{"name":"dq.WatermarksLateArrivalDelayMs"},{"name":"dq.WatermarksMode"},{"name":"dq.WorkerFilter"},{"name":"dq.WorkersPerOperation"},{"name":"AllowDotInAlias"},{"name":"AllowUnnamedColumns"},{"name":"AnsiCurrentRow"},{"name":"AnsiImplicitCrossJoin"},{"name":"AnsiInForEmptyOrNullableItemsCollections"},{"name":"AnsiLike"},{"name":"AnsiOptionalAs"},{"name":"AnsiRankForNullableKeys"},{"name":"AutoCommit"},{"name":"BlockEngine"},{"name":"BlockEngineEnable"},{"name":"BlockEngineForce"},{"name":"BogousStarInGroupByOverJoin"},{"name":"CheckedOps"},{"name":"ClassicDivision"},{"name":"CoalesceJoinKeysOnQualifiedAll"},{"name":"CompactGroupBy"},{"name":"CompactNamedExprs"},{"name":"CostBasedOptimizer"},{"name":"DataWatermarks"},{"name":"DirectRead"},{"name":"DisableAnsiCurrentRow"},{"name":"DisableAnsiImplicitCrossJoin"},{"name":"DisableAnsiInForEmptyOrNullableItemsCollections"},{"name":"DisableAnsiLike"},{"name":"DisableAnsiOptionalAs"},{"name":"DisableAnsiRankForNullableKeys"},{"name":"DisableBlockEngineEnable"},{"name":"DisableBlockEngineForce"},{"name":"DisableBogousStarInGroupByOverJoin"},{"name":"DisableCoalesceJoinKeysOnQualifiedAll"},{"name":"DisableCompactGroupBy"},{"name":"DisableCompactNamedExprs"},{"name":"DisableDistinctOverKeys"},{"name":"DisableDistinctOverWindow"},{"name":"DisableDqEngineEnable"},{"name":"DisableDqEngineForce"},{"name":"DisableEmitAggApply"},{"name":"DisableEmitStartsWith"},{"name":"DisableEmitTableSource"},{"name":"DisableEmitUnionMerge"},{"name":"DisableFilterPushdownOverJoinOptionalSide"},{"name":"DisableFlexibleTypes"},{"name":"DisableGroupByExprAfterWhere"},{"name":"DisableJsonQueryReturnsJsonDocument"},{"name":"DisableOrderedColumns"},{"name":"DisablePullUpFlatMapOverJoin"},{"name":"DisableRegexUseRe2"},{"name":"DisableRotateJoinTree"},{"name":"DisableSeqMode"},{"name":"DisableSimpleColumns"},{"name":"DisableStrictJoinKeyTypes"},{"name":"DisableUnicodeLiterals"},{"name":"DisableUnorderedResult"},{"name":"DisableUnorderedSubqueries"},{"name":"DisableUseBlocks"},{"name":"DisableValidateUnusedExprs"},{"name":"DisableWarnOnAnsiAliasShadowing"},{"name":"DisableWarnUntypedStringLiterals"},{"name":"DiscoveryMode"},{"name":"DistinctOverKeys"},{"name":"DistinctOverWindow"},{"name":"DqEngine"},{"name":"DqEngineEnable"},{"name":"DqEngineForce"},{"name":"EmitAggApply"},{"name":"EmitStartsWith"},{"name":"EmitTableSource"},{"name":"EmitUnionMerge"},{"name":"EnableSystemColumns"},{"name":"Engine"},{"name":"ErrorMsg"},{"name":"FeatureR010"},{"name":"File"},{"name":"FileOption"},{"name":"FilterPushdownOverJoinOptionalSide"},{"name":"FlexibleTypes"},{"name":"Folder"},{"name":"Greetings"},{"name":"GroupByCubeLimit"},{"name":"GroupByExprAfterWhere"},{"name":"GroupByLimit"},{"name":"JsonQueryReturnsJsonDocument"},{"name":"Library"},{"name":"OrderedColumns"},{"name":"OverrideLibrary"},{"name":"Package"},{"name":"PackageVersion"},{"name":"PathPrefix"},{"name":"PositionalUnionAll"},{"name":"PqReadBy"},{"name":"PullUpFlatMapOverJoin"},{"name":"RefSelect"},{"name":"RegexUseRe2"},{"name":"ResultRowsLimit"},{"name":"ResultSizeLimit"},{"name":"RotateJoinTree"},{"name":"RuntimeLogLevel"},{"name":"SampleSelect"},{"name":"SeqMode"},{"name":"SimpleColumns"},{"name":"StrictJoinKeyTypes"},{"name":"Udf"},{"name":"UnicodeLiterals"},{"name":"UnorderedResult"},{"name":"UnorderedSubqueries"},{"name":"UseBlocks"},{"name":"UseTablePrefixForEach"},{"name":"ValidateUnusedExprs"},{"name":"WarnOnAnsiAliasShadowing"},{"name":"WarnUnnamedColumns"},{"name":"WarnUntypedStringLiterals"},{"name":"Warning"},{"name":"WarningMsg"},{"name":"yson.AutoConvert"},{"name":"yson.CastToString"},{"name":"yson.DisableCastToString"},{"name":"yson.DisableStrict"},{"name":"yson.Strict"}] diff --git a/yql/essentials/minikql/comp_nodes/mkql_condense.cpp b/yql/essentials/minikql/comp_nodes/mkql_condense.cpp index ced3a95d5ab..04181778482 100644 --- a/yql/essentials/minikql/comp_nodes/mkql_condense.cpp +++ b/yql/essentials/minikql/comp_nodes/mkql_condense.cpp @@ -30,7 +30,7 @@ public: NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { if (state.IsFinish()) { - return static_cast<const NUdf::TUnboxedValuePod&>(state); + return state; } if (state.IsInvalid()) { diff --git a/yql/essentials/minikql/comp_nodes/mkql_condense1.cpp b/yql/essentials/minikql/comp_nodes/mkql_condense1.cpp index 850a7bff042..b2cda5a5872 100644 --- a/yql/essentials/minikql/comp_nodes/mkql_condense1.cpp +++ b/yql/essentials/minikql/comp_nodes/mkql_condense1.cpp @@ -30,7 +30,7 @@ public: NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { if (state.IsFinish()) { - return static_cast<const NUdf::TUnboxedValuePod&>(state); + return state; } else if (state.HasValue()) { if constexpr (UseCtx) { CleanupCurrentContext(); diff --git a/yql/essentials/minikql/comp_nodes/mkql_flatmap.cpp b/yql/essentials/minikql/comp_nodes/mkql_flatmap.cpp index 923974e9e7e..69c42c7afbd 100644 --- a/yql/essentials/minikql/comp_nodes/mkql_flatmap.cpp +++ b/yql/essentials/minikql/comp_nodes/mkql_flatmap.cpp @@ -219,7 +219,7 @@ public: } if (state.IsFinish()) { - return NUdf::TUnboxedValuePod::MakeFinish(); + return state; } while (true) { diff --git a/yql/essentials/minikql/comp_nodes/mkql_multimap.cpp b/yql/essentials/minikql/comp_nodes/mkql_multimap.cpp index 1a0639b9e9b..5e33c133c31 100644 --- a/yql/essentials/minikql/comp_nodes/mkql_multimap.cpp +++ b/yql/essentials/minikql/comp_nodes/mkql_multimap.cpp @@ -21,8 +21,9 @@ public: {} NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { - if (state.IsFinish()) - return NUdf::TUnboxedValuePod::MakeFinish(); + if (state.IsFinish()) { + return state; + } const auto pos = state.IsInvalid() ? 0ULL : state.Get<ui64>(); if (!pos) { @@ -412,8 +413,9 @@ public: {} NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { - if (state.IsFinish()) - return NUdf::TUnboxedValuePod::MakeFinish(); + if (state.IsFinish()) { + return state; + } const auto pos = state.IsInvalid() ? 0ULL : state.Get<ui64>(); if (!pos) { diff --git a/yql/essentials/minikql/comp_nodes/mkql_squeeze_to_list.cpp b/yql/essentials/minikql/comp_nodes/mkql_squeeze_to_list.cpp index 548eb0937df..55bd62b112b 100644 --- a/yql/essentials/minikql/comp_nodes/mkql_squeeze_to_list.cpp +++ b/yql/essentials/minikql/comp_nodes/mkql_squeeze_to_list.cpp @@ -47,7 +47,7 @@ public: NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { if (state.IsFinish()) { - return NUdf::TUnboxedValuePod::MakeFinish(); + return state; } else if (state.IsInvalid()) { MakeState(ctx, Limit->GetValue(ctx).GetOrDefault(std::numeric_limits<ui64>::max()), state); } diff --git a/yql/essentials/minikql/comp_nodes/mkql_todict.cpp b/yql/essentials/minikql/comp_nodes/mkql_todict.cpp index 738f88232a9..0ca6c68166a 100644 --- a/yql/essentials/minikql/comp_nodes/mkql_todict.cpp +++ b/yql/essentials/minikql/comp_nodes/mkql_todict.cpp @@ -985,7 +985,7 @@ public: NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { if (state.IsFinish()) { - return state.Release(); + return state; } else if (state.IsInvalid()) { MakeState(ctx, state); } @@ -1162,7 +1162,7 @@ public: NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { if (state.IsFinish()) { - return state.Release(); + return state; } else if (state.IsInvalid()) { MakeState(ctx, state); } diff --git a/yql/essentials/minikql/comp_nodes/ut/mkql_todict_ut.cpp b/yql/essentials/minikql/comp_nodes/ut/mkql_todict_ut.cpp index abb4f5d85e9..f09b9529d26 100644 --- a/yql/essentials/minikql/comp_nodes/ut/mkql_todict_ut.cpp +++ b/yql/essentials/minikql/comp_nodes/ut/mkql_todict_ut.cpp @@ -147,6 +147,10 @@ Y_UNIT_TEST_SUITE(TMiniKQLToDictTest) { status = res.Fetch(v); UNIT_ASSERT_VALUES_EQUAL(NUdf::EFetchStatus::Finish, status); + // XXX: Check whether the internal state is not released + // and the sentinel is still set (see more info in YQL-19866). + status = res.Fetch(v); + UNIT_ASSERT_VALUES_EQUAL(NUdf::EFetchStatus::Finish, status); }; for (auto stream : {true, false}) { @@ -201,6 +205,10 @@ Y_UNIT_TEST_SUITE(TMiniKQLToDictTest) { status = res.Fetch(v); UNIT_ASSERT_VALUES_EQUAL(NUdf::EFetchStatus::Finish, status); + // XXX: Check whether the internal state is not released + // and the sentinel is still set (see more info in YQL-19866). + status = res.Fetch(v); + UNIT_ASSERT_VALUES_EQUAL(NUdf::EFetchStatus::Finish, status); }; for (auto hashed : {true, false}) { diff --git a/yql/essentials/sql/v1/context.cpp b/yql/essentials/sql/v1/context.cpp index 1a0a1f4b18d..7f3d5433b96 100644 --- a/yql/essentials/sql/v1/context.cpp +++ b/yql/essentials/sql/v1/context.cpp @@ -69,6 +69,7 @@ THashMap<TStringBuf, TPragmaField> CTX_PRAGMA_FIELDS = { {"EmitUnionMerge", &TContext::EmitUnionMerge}, {"SeqMode", &TContext::SeqMode}, {"DistinctOverKeys", &TContext::DistinctOverKeys}, + {"GroupByExprAfterWhere", &TContext::GroupByExprAfterWhere}, }; typedef TMaybe<bool> TContext::*TPragmaMaybeField; @@ -104,6 +105,10 @@ TContext::TContext(const TLexers& lexers, const TParsers& parsers, , WarningPolicy(settings.IsReplay) , BlockEngineEnable(Settings.BlockDefaultAuto->Allow()) { + if (settings.LangVer >= MakeLangVersion(2025, 2)) { + GroupByExprAfterWhere = true; + } + for (auto lib : settings.Libraries) { Libraries.emplace(lib, TLibraryStuff()); } diff --git a/yql/essentials/sql/v1/context.h b/yql/essentials/sql/v1/context.h index 3bdfa1ceab4..e3f70f7515e 100644 --- a/yql/essentials/sql/v1/context.h +++ b/yql/essentials/sql/v1/context.h @@ -373,6 +373,7 @@ namespace NSQLTranslationV1 { bool DistinctOverWindow = false; bool SeqMode = false; bool DistinctOverKeys = false; + bool GroupByExprAfterWhere = false; bool EmitUnionMerge = false; TVector<size_t> ForAllStatementsParts; diff --git a/yql/essentials/sql/v1/select.cpp b/yql/essentials/sql/v1/select.cpp index c4290f3268f..cb72dc95f5f 100644 --- a/yql/essentials/sql/v1/select.cpp +++ b/yql/essentials/sql/v1/select.cpp @@ -1744,6 +1744,11 @@ public: if (Flatten) { block = L(block, Y("let", "core", Y(ordered ? "OrderedFlatMap" : "FlatMap", "core", BuildLambda(Pos, Y("row"), Flatten, "res")))); } + if (ctx.GroupByExprAfterWhere) { + if (auto filter = Source->BuildFilter(ctx, "core"); filter) { + block = L(block, Y("let", "core", filter)); + } + } if (PreaggregatedMap) { block = L(block, Y("let", "core", PreaggregatedMap)); if (Source->IsCompositeSource() && !Columns.QualifiedAll) { @@ -1752,9 +1757,10 @@ public: } else if (Source->IsCompositeSource() && !Columns.QualifiedAll) { block = L(block, Y("let", "origcore", "core")); } - auto filter = Source->BuildFilter(ctx, "core"); - if (filter) { - block = L(block, Y("let", "core", filter)); + if (!ctx.GroupByExprAfterWhere) { + if (auto filter = Source->BuildFilter(ctx, "core"); filter) { + block = L(block, Y("let", "core", filter)); + } } if (Aggregate) { block = L(block, Y("let", "core", Aggregate)); diff --git a/yql/essentials/sql/v1/sql_query.cpp b/yql/essentials/sql/v1/sql_query.cpp index b60a338d02f..d64a7bf8d61 100644 --- a/yql/essentials/sql/v1/sql_query.cpp +++ b/yql/essentials/sql/v1/sql_query.cpp @@ -3393,6 +3393,12 @@ TNodePtr TSqlQuery::PragmaStatement(const TRule_pragma_stmt& stmt, bool& success } else if (normalizedPragma == "disabledistinctoverkeys") { Ctx.DistinctOverKeys = false; Ctx.IncrementMonCounter("sql_pragma", "DisableDistinctOverKeys"); + } else if (normalizedPragma == "groupbyexprafterwhere") { + Ctx.GroupByExprAfterWhere = true; + Ctx.IncrementMonCounter("sql_pragma", "GroupByExprAfterWhere"); + } else if (normalizedPragma == "disablegroupbyexprafterwhere") { + Ctx.GroupByExprAfterWhere = false; + Ctx.IncrementMonCounter("sql_pragma", "DisableGroupByExprAfterWhere"); } else if (normalizedPragma == "engine") { Ctx.IncrementMonCounter("sql_pragma", "Engine"); diff --git a/yql/essentials/tests/common/test_framework/test_file_common.py b/yql/essentials/tests/common/test_framework/test_file_common.py index 240182e0056..9c91f8a1671 100644 --- a/yql/essentials/tests/common/test_framework/test_file_common.py +++ b/yql/essentials/tests/common/test_framework/test_file_common.py @@ -88,7 +88,7 @@ def get_sql_query(provider, suite, case, config, data_path=None, template='.sql' def run_file_no_cache(provider, suite, case, cfg, config, yql_http_file_server, yqlrun_binary=None, extra_args=[], force_blocks=False, allow_llvm=True, data_path=None, - run_sql=True, cfg_postprocess=None): + run_sql=True, cfg_postprocess=None, langver=None): check_provider(provider, config) sql_query = get_sql_query(provider, suite, case, config, data_path, template='.sql' if run_sql else '.yqls') @@ -119,7 +119,8 @@ def run_file_no_cache(provider, suite, case, cfg, config, yql_http_file_server, gateway_config=get_gateways_config(http_files, yql_http_file_server, force_blocks=force_blocks, is_hybrid=is_hybrid(provider), allow_llvm=allow_llvm, postprocess_func=cfg_postprocess), extra_args=extra_args, - udfs_dir=yql_binary_path('yql/essentials/tests/common/test_framework/udfs_deps') + udfs_dir=yql_binary_path('yql/essentials/tests/common/test_framework/udfs_deps'), + langver=langver ) res, tables_res = execute( @@ -156,12 +157,14 @@ def run_file_no_cache(provider, suite, case, cfg, config, yql_http_file_server, def run_file(provider, suite, case, cfg, config, yql_http_file_server, yqlrun_binary=None, - extra_args=[], force_blocks=False, allow_llvm=True, data_path=None, run_sql=True, cfg_postprocess=None): + extra_args=[], force_blocks=False, allow_llvm=True, data_path=None, run_sql=True, + cfg_postprocess=None, langver=None): if (suite, case, cfg) not in run_file.cache: run_file.cache[(suite, case, cfg)] = \ run_file_no_cache(provider, suite, case, cfg, config, yql_http_file_server, yqlrun_binary, extra_args, force_blocks=force_blocks, allow_llvm=allow_llvm, - data_path=data_path, run_sql=run_sql, cfg_postprocess=cfg_postprocess) + data_path=data_path, run_sql=run_sql, cfg_postprocess=cfg_postprocess, + langver=langver) return run_file.cache[(suite, case, cfg)] diff --git a/yql/essentials/tests/common/test_framework/test_utils.py b/yql/essentials/tests/common/test_framework/test_utils.py index 0865245fd9b..e760e19a342 100644 --- a/yql/essentials/tests/common/test_framework/test_utils.py +++ b/yql/essentials/tests/common/test_framework/test_utils.py @@ -144,6 +144,7 @@ def validate_cfg(result): "yt_file", "os", "param", + "langver", ), "Unknown command in .cfg: %s" % (r[0]) diff --git a/yql/essentials/tests/common/test_framework/yql_utils.py b/yql/essentials/tests/common/test_framework/yql_utils.py index 3e4a4afa3fe..2d59a1fa13f 100644 --- a/yql/essentials/tests/common/test_framework/yql_utils.py +++ b/yql/essentials/tests/common/test_framework/yql_utils.py @@ -496,6 +496,13 @@ def is_xfail(cfg): return False +def get_langver(cfg): + for item in cfg: + if item[0] == 'langver': + return item[1] + return None + + def is_skip_forceblocks(cfg): for item in cfg: if item[0] == 'skip_forceblocks': diff --git a/yql/essentials/tests/common/test_framework/yqlrun.py b/yql/essentials/tests/common/test_framework/yqlrun.py index e37bf7c38d3..49086ba7e88 100644 --- a/yql/essentials/tests/common/test_framework/yqlrun.py +++ b/yql/essentials/tests/common/test_framework/yqlrun.py @@ -25,7 +25,8 @@ FIX_DIR_PREFIXES = { class YQLRun(object): - def __init__(self, udfs_dir=None, prov='yt', use_sql2yql=False, keep_temp=True, binary=None, gateway_config=None, fs_config=None, extra_args=[], cfg_dir=None, support_udfs=True): + def __init__(self, udfs_dir=None, prov='yt', use_sql2yql=False, keep_temp=True, binary=None, gateway_config=None, + fs_config=None, extra_args=[], cfg_dir=None, support_udfs=True, langver=None): if binary is None: self.yqlrun_binary = yql_utils.yql_binary_path(os.getenv('YQL_YQLRUN_PATH') or 'yql/tools/yqlrun/yqlrun') else: @@ -80,6 +81,8 @@ class YQLRun(object): flags = yql_utils.get_param('SQL_FLAGS').split(',') self.gateway_config.SqlCore.TranslationFlags.extend(flags) + self.langver = langver + def yql_exec(self, program=None, program_file=None, files=None, urls=None, run_sql=False, verbose=False, check_error=True, tables=None, pretty_plan=True, wait=True, parameters={}, extra_env={}, require_udf_resolver=False, scan_udfs=True): @@ -173,6 +176,9 @@ class YQLRun(object): if ansi_lexer: cmd += '--ansi-lexer ' + if self.langver is not None: + cmd += '--langver=%s ' % (self.langver,) + if self.keep_temp and prov != 'pure': cmd += '--keep-temp ' diff --git a/yql/essentials/tests/s-expressions/minirun/pure.py b/yql/essentials/tests/s-expressions/minirun/pure.py index 38318d9fb8a..576230fd00b 100644 --- a/yql/essentials/tests/s-expressions/minirun/pure.py +++ b/yql/essentials/tests/s-expressions/minirun/pure.py @@ -9,12 +9,13 @@ from yql_utils import execute, get_tables, get_files, get_http_files, \ KSV_ATTR, yql_binary_path, is_xfail, is_canonize_peephole, is_peephole_use_blocks, is_canonize_lineage, \ is_skip_forceblocks, get_param, normalize_source_code_path, replace_vals, get_gateway_cfg_suffix, \ do_custom_query_check, stable_result_file, stable_table_file, is_with_final_result_issues, \ - normalize_result + normalize_result, get_langver from yqlrun import YQLRun from test_utils import get_config, get_parameters_json from test_file_common import run_file, run_file_no_cache, get_gateways_config, get_sql_query +DEFAULT_LANG_VER = '2025.01' ASTDIFF_PATH = yql_binary_path('yql/essentials/tools/astdiff/astdiff') MINIRUN_PATH = yql_binary_path('yql/essentials/tools/minirun/minirun') DATA_PATH = yatest.common.source_path('yql/essentials/tests/s-expressions/suites') @@ -25,6 +26,9 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): pytest.skip('non-trivial gateways.conf') config = get_config(suite, case, cfg, data_path=DATA_PATH) + langver = get_langver(config) + if langver is None: + langver = DEFAULT_LANG_VER xfail = is_xfail(config) if xfail and what != 'Results': @@ -38,7 +42,8 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): if is_with_final_result_issues(config): extra_final_args += ['--with-final-issues'] (res, tables_res) = run_file('pure', suite, case, cfg, config, yql_http_file_server, MINIRUN_PATH, - extra_args=extra_final_args, allow_llvm=False, data_path=DATA_PATH, run_sql=False) + extra_args=extra_final_args, allow_llvm=False, data_path=DATA_PATH, + run_sql=False, langver=langver) to_canonize = [] assert not tables_res @@ -72,7 +77,8 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): keep_temp=False, gateway_config=get_gateways_config(http_files, yql_http_file_server, allow_llvm=is_llvm), udfs_dir=yql_binary_path('yql/essentials/tests/common/test_framework/udfs_deps'), - binary=MINIRUN_PATH + binary=MINIRUN_PATH, + langver=langver ) opt_res, opt_tables_res = execute( diff --git a/yql/essentials/tests/sql/minirun/part5/canondata/result.json b/yql/essentials/tests/sql/minirun/part5/canondata/result.json index 579677b9314..ec17a6eec0b 100644 --- a/yql/essentials/tests/sql/minirun/part5/canondata/result.json +++ b/yql/essentials/tests/sql/minirun/part5/canondata/result.json @@ -237,6 +237,20 @@ "uri": "https://{canondata_backend}/1871002/ab54d2c5acdb4e70fca2cf294e5ea9c225baab0c/resource.tar.gz#test.test_aggr_factory-transform_output-default.txt-Results_/results.txt" } ], + "test.test[aggregate-group_by_expr_after_where_ver--Debug]": [ + { + "checksum": "d0da2a1dc674dc86930994f910eb1c2b", + "size": 253, + "uri": "https://{canondata_backend}/1130705/42827d049e4963da219fd249860b757d672765ec/resource.tar.gz#test.test_aggregate-group_by_expr_after_where_ver--Debug_/opt.yql" + } + ], + "test.test[aggregate-group_by_expr_after_where_ver--Results]": [ + { + "checksum": "1578f70c73d3ed10cf371be4ce4e70ec", + "size": 690, + "uri": "https://{canondata_backend}/1130705/42827d049e4963da219fd249860b757d672765ec/resource.tar.gz#test.test_aggregate-group_by_expr_after_where_ver--Results_/results.txt" + } + ], "test.test[aggregate-hopping-default.txt-Debug]": [ { "checksum": "bb23bcbf33c639f7faf92653dd668f3c", @@ -1044,6 +1058,20 @@ "uri": "https://{canondata_backend}/1936273/19f08c34eba9366d29ee0ffb8eb99e637c34fd97/resource.tar.gz#test.test_join-convert_check_key_mem2-default.txt-Results_/results.txt" } ], + "test.test[join-prune_keys-default.txt-Debug]": [ + { + "checksum": "a706e6bd4285c96ad5120d701109fc82", + "size": 3669, + "uri": "https://{canondata_backend}/1871102/392832e505c55eb371c9d3241b89c96b5a837c8f/resource.tar.gz#test.test_join-prune_keys-default.txt-Debug_/opt.yql" + } + ], + "test.test[join-prune_keys-default.txt-Results]": [ + { + "checksum": "af076a3334031b8cdaa6969f064a4616", + "size": 18160, + "uri": "https://{canondata_backend}/1130705/620da5a4f19baef17c32a4b3c699ec3c3091ada5/resource.tar.gz#test.test_join-prune_keys-default.txt-Results_/results.txt" + } + ], "test.test[json-json_exists/common_syntax-default.txt-Debug]": [ { "checksum": "1559e7b19e1d1827f8f1ea62929effb1", diff --git a/yql/essentials/tests/sql/minirun/part6/canondata/result.json b/yql/essentials/tests/sql/minirun/part6/canondata/result.json index 26ca736cfcf..60e9fee9459 100644 --- a/yql/essentials/tests/sql/minirun/part6/canondata/result.json +++ b/yql/essentials/tests/sql/minirun/part6/canondata/result.json @@ -195,6 +195,20 @@ "uri": "https://{canondata_backend}/1936273/614fe8dff439fd011c07c47361f2a1d0d854297f/resource.tar.gz#test.test_aggregate-distinct_over_keys-default.txt-Results_/results.txt" } ], + "test.test[aggregate-group_by_expr_after_where-default.txt-Debug]": [ + { + "checksum": "d0da2a1dc674dc86930994f910eb1c2b", + "size": 253, + "uri": "https://{canondata_backend}/1925842/9a344928381729abc8381a7e8ada7e10e2ba51fe/resource.tar.gz#test.test_aggregate-group_by_expr_after_where-default.txt-Debug_/opt.yql" + } + ], + "test.test[aggregate-group_by_expr_after_where-default.txt-Results]": [ + { + "checksum": "1578f70c73d3ed10cf371be4ce4e70ec", + "size": 690, + "uri": "https://{canondata_backend}/1925842/9a344928381729abc8381a7e8ada7e10e2ba51fe/resource.tar.gz#test.test_aggregate-group_by_expr_after_where-default.txt-Results_/results.txt" + } + ], "test.test[ansi_idents-escaping-default.txt-Debug]": [ { "checksum": "13dd30dd58fd993aa21441bec427f12b", diff --git a/yql/essentials/tests/sql/minirun/part7/canondata/result.json b/yql/essentials/tests/sql/minirun/part7/canondata/result.json index e8c9926ae75..8083b15dac8 100644 --- a/yql/essentials/tests/sql/minirun/part7/canondata/result.json +++ b/yql/essentials/tests/sql/minirun/part7/canondata/result.json @@ -1232,16 +1232,16 @@ ], "test.test[select-prune_keys-default.txt-Debug]": [ { - "checksum": "95e58e469ce10fce0d0d5e55c0cf3baf", - "size": 3292, - "uri": "https://{canondata_backend}/1931696/04008bc01ad4f562f8e03ad2bc296f7a64a78489/resource.tar.gz#test.test_select-prune_keys-default.txt-Debug_/opt.yql" + "checksum": "1dc9319da10b4c1343b6ceb5d0abb5b7", + "size": 3788, + "uri": "https://{canondata_backend}/212715/392992a262a39acb9e6e49f104e5800a3e731eb6/resource.tar.gz#test.test_select-prune_keys-default.txt-Debug_/opt.yql" } ], "test.test[select-prune_keys-default.txt-Results]": [ { - "checksum": "f53b6976b6a5e1ee5e1d331c25963930", - "size": 27670, - "uri": "https://{canondata_backend}/1931696/04008bc01ad4f562f8e03ad2bc296f7a64a78489/resource.tar.gz#test.test_select-prune_keys-default.txt-Results_/results.txt" + "checksum": "46c226ab1c00ed4a9a5b57c6639f15b2", + "size": 30306, + "uri": "https://{canondata_backend}/212715/392992a262a39acb9e6e49f104e5800a3e731eb6/resource.tar.gz#test.test_select-prune_keys-default.txt-Results_/results.txt" } ], "test.test[union-union_positional_mix-default.txt-Debug]": [ diff --git a/yql/essentials/tests/sql/minirun/pure.py b/yql/essentials/tests/sql/minirun/pure.py index c3a78adf354..3ec17bb2fa0 100644 --- a/yql/essentials/tests/sql/minirun/pure.py +++ b/yql/essentials/tests/sql/minirun/pure.py @@ -9,12 +9,13 @@ from yql_utils import execute, get_tables, get_files, get_http_files, \ KSV_ATTR, yql_binary_path, is_xfail, is_canonize_peephole, is_peephole_use_blocks, is_canonize_lineage, \ is_skip_forceblocks, get_param, normalize_source_code_path, replace_vals, get_gateway_cfg_suffix, \ do_custom_query_check, stable_result_file, stable_table_file, is_with_final_result_issues, \ - normalize_result + normalize_result, get_langver from yqlrun import YQLRun from test_utils import get_config, get_parameters_json from test_file_common import run_file, run_file_no_cache, get_gateways_config, get_sql_query +DEFAULT_LANG_VER = '2025.01' DATA_PATH = yatest.common.source_path('yql/essentials/tests/sql/suites') ASTDIFF_PATH = yql_binary_path('yql/essentials/tools/astdiff/astdiff') MINIRUN_PATH = yql_binary_path('yql/essentials/tools/minirun/minirun') @@ -39,6 +40,10 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): config = get_config(suite, case, cfg, data_path = DATA_PATH) + langver = get_langver(config) + if langver is None: + langver = DEFAULT_LANG_VER + xfail = is_xfail(config) if xfail and what != 'Results': pytest.skip('xfail is not supported in this mode') @@ -51,7 +56,8 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): if is_with_final_result_issues(config): extra_final_args += ['--with-final-issues'] (res, tables_res) = run_file('pure', suite, case, cfg, config, yql_http_file_server, MINIRUN_PATH, - extra_args=extra_final_args, allow_llvm=False, data_path=DATA_PATH) + extra_args=extra_final_args, allow_llvm=False, data_path=DATA_PATH, + langver=langver) to_canonize = [] assert xfail or os.path.exists(res.results_file) @@ -67,7 +73,8 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): force_blocks = is_peephole_use_blocks(config) (res, tables_res) = run_file_no_cache('pure', suite, case, cfg, config, yql_http_file_server, force_blocks=force_blocks, extra_args=['--peephole'], - data_path=DATA_PATH, yqlrun_binary=MINIRUN_PATH) + data_path=DATA_PATH, yqlrun_binary=MINIRUN_PATH, + langver=langver) return [yatest.common.canonical_file(res.opt_file, diff_tool=ASTDIFF_PATH)] if what == 'Results': @@ -100,7 +107,8 @@ def run_test(suite, case, cfg, tmpdir, what, yql_http_file_server): keep_temp=False, gateway_config=get_gateways_config(http_files, yql_http_file_server, allow_llvm=is_llvm, force_blocks=is_blocks), udfs_dir=yql_binary_path('yql/essentials/tests/common/test_framework/udfs_deps'), - binary=MINIRUN_PATH + binary=MINIRUN_PATH, + langver=langver ) opt_res, opt_tables_res = execute( diff --git a/yql/essentials/tests/sql/sql2yql/canondata/result.json b/yql/essentials/tests/sql/sql2yql/canondata/result.json index 4cb7e6cc7af..1022b92f85c 100644 --- a/yql/essentials/tests/sql/sql2yql/canondata/result.json +++ b/yql/essentials/tests/sql/sql2yql/canondata/result.json @@ -937,6 +937,20 @@ "uri": "https://{canondata_backend}/1936273/e22f8123b51c2802f50d5a8d4626267f2f28e9ab/resource.tar.gz#test_sql2yql.test_aggregate-distinct_over_keys_/sql.yql" } ], + "test_sql2yql.test[aggregate-group_by_expr_after_where]": [ + { + "checksum": "4799645fef77850f5f5f07de2d1b8bc2", + "size": 1762, + "uri": "https://{canondata_backend}/1899731/49525280cc90ece19469c3347e616ee12710ec2c/resource.tar.gz#test_sql2yql.test_aggregate-group_by_expr_after_where_/sql.yql" + } + ], + "test_sql2yql.test[aggregate-group_by_expr_after_where_ver]": [ + { + "checksum": "d7d81aab522ef18bb07e619c426594c5", + "size": 1762, + "uri": "https://{canondata_backend}/1597364/99b2cf59a9975dbc2994ead01aa9dcbd784b5279/resource.tar.gz#test_sql2yql.test_aggregate-group_by_expr_after_where_ver_/sql.yql" + } + ], "test_sql2yql.test[aggregate-group_by_rollup_rename]": [ { "checksum": "bc5b27508587d82ba3e9d0a752d25dcc", @@ -3933,6 +3947,13 @@ "uri": "https://{canondata_backend}/1942173/99e88108149e222741552e7e6cddef041d6a2846/resource.tar.gz#test_sql2yql.test_join-left_join_with_self_aggr_/sql.yql" } ], + "test_sql2yql.test[join-prune_keys]": [ + { + "checksum": "a04490e5ef3a567ca13cb1ed88272bd8", + "size": 22535, + "uri": "https://{canondata_backend}/1599023/2a161150407124ac83c2566a6542b19a53abfccb/resource.tar.gz#test_sql2yql.test_join-prune_keys_/sql.yql" + } + ], "test_sql2yql.test[join-yql-19192]": [ { "checksum": "fffdf1cbb40643da9daf9bdf3edec121", @@ -7022,9 +7043,9 @@ ], "test_sql2yql.test[select-prune_keys]": [ { - "checksum": "55346f77548ef19f9a09d2f1d3f6f466", - "size": 17765, - "uri": "https://{canondata_backend}/1871182/906a4c4e540bb8746f8d7595500d4d1c9f664846/resource.tar.gz#test_sql2yql.test_select-prune_keys_/sql.yql" + "checksum": "f7da5706622461ab177712e6c348c61b", + "size": 19536, + "uri": "https://{canondata_backend}/1814674/c8d78993e8e9976f1e3fae2197140afe33195365/resource.tar.gz#test_sql2yql.test_select-prune_keys_/sql.yql" } ], "test_sql2yql.test[select-result_label]": [ @@ -8103,6 +8124,16 @@ "uri": "file://test_sql_format.test_aggregate-distinct_over_keys_/formatted.sql" } ], + "test_sql_format.test[aggregate-group_by_expr_after_where]": [ + { + "uri": "file://test_sql_format.test_aggregate-group_by_expr_after_where_/formatted.sql" + } + ], + "test_sql_format.test[aggregate-group_by_expr_after_where_ver]": [ + { + "uri": "file://test_sql_format.test_aggregate-group_by_expr_after_where_ver_/formatted.sql" + } + ], "test_sql_format.test[aggregate-group_by_rollup_rename]": [ { "uri": "file://test_sql_format.test_aggregate-group_by_rollup_rename_/formatted.sql" @@ -10243,6 +10274,11 @@ "uri": "file://test_sql_format.test_join-left_join_with_self_aggr_/formatted.sql" } ], + "test_sql_format.test[join-prune_keys]": [ + { + "uri": "file://test_sql_format.test_join-prune_keys_/formatted.sql" + } + ], "test_sql_format.test[join-yql-19192]": [ { "uri": "file://test_sql_format.test_join-yql-19192_/formatted.sql" diff --git a/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_aggregate-group_by_expr_after_where_/formatted.sql b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_aggregate-group_by_expr_after_where_/formatted.sql new file mode 100644 index 00000000000..cb490bcb76c --- /dev/null +++ b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_aggregate-group_by_expr_after_where_/formatted.sql @@ -0,0 +1,13 @@ +PRAGMA GroupByExprAfterWhere; + +SELECT + x +FROM ( + SELECT + 1 AS x +) +WHERE + x == 1 +GROUP BY + -x AS x +; diff --git a/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_aggregate-group_by_expr_after_where_ver_/formatted.sql b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_aggregate-group_by_expr_after_where_ver_/formatted.sql new file mode 100644 index 00000000000..20d25d53d24 --- /dev/null +++ b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_aggregate-group_by_expr_after_where_ver_/formatted.sql @@ -0,0 +1,11 @@ +SELECT + x +FROM ( + SELECT + 1 AS x +) +WHERE + x == 1 +GROUP BY + -x AS x +; diff --git a/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_join-prune_keys_/formatted.sql b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_join-prune_keys_/formatted.sql new file mode 100644 index 00000000000..0f53bcbe441 --- /dev/null +++ b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_join-prune_keys_/formatted.sql @@ -0,0 +1,238 @@ +PRAGMA config.flags('OptimizerFlags', 'EmitPruneKeys'); + +$a = ( + SELECT + * + FROM + as_table([ + <|x: 1, t: 1|>, + <|x: 1, t: 1|>, + <|x: 1, t: 2|>, + <|x: 3, t: 1|>, + <|x: 3, t: 4|>, + <|x: 3, t: 2|>, + ]) +); + +$b = ( + SELECT + * + FROM + as_table([ + <|x: 1, y: 1|>, + <|x: 1, y: 2|>, + <|x: 1, y: 3|>, + <|x: 1, y: 3|>, + <|x: 2, y: 3|>, + <|x: 2, y: 4|>, + ]) +); + +$c = ( + SELECT + * + FROM + as_table([ + <|x: 1|>, + <|x: 1|>, + <|x: 1|>, + <|x: 1|>, + <|x: 2|>, + <|x: 2|>, + ]) +); + +-- PruneKeys +SELECT + a.* +FROM + $a AS a +WHERE + a.x IN ( + SELECT + x + FROM + $b + ) +; -- PruneKeys + +SELECT + a.* +FROM + $a AS a +WHERE + a.x IN ( + SELECT + /*+ distinct(x) */ x + FROM + $b + ) +; -- nothing + +SELECT + a.* +FROM + $a AS a +WHERE + a.x IN ( + SELECT + x + FROM + $c + ) +; -- PruneKeys + +SELECT + a.* +FROM + $a AS a +LEFT SEMI JOIN + $b AS b +ON + a.x == b.x +; -- PruneKeys(b) + +SELECT + a.* +FROM + $b AS b +RIGHT SEMI JOIN + $a AS a +ON + b.x == a.x +; -- PruneKeys(b) + +SELECT + a.x, + a.t, + b.x +FROM ANY + $a AS a +JOIN + $b AS b +ON + a.x == b.x +; -- PruneKeys(a) + +SELECT + a.x, + a.t, + b.x +FROM + $a AS a +JOIN ANY + $b AS b +ON + a.x == b.x +; -- PruneKeys(b) + +$a_sorted = ( + SELECT + * + FROM + $a + ASSUME ORDER BY + x +); + +$b_sorted = ( + SELECT + * + FROM + $b + ASSUME ORDER BY + x +); + +$c_sorted = ( + SELECT + * + FROM + $c + ASSUME ORDER BY + x +); + +-- PruneAdjacentKeys +SELECT + a.* +FROM + $a AS a +WHERE + a.x IN ( + SELECT + x + FROM + $b_sorted + ) +; -- PruneAdjacentKeys + +SELECT + a.* +FROM + $a AS a +WHERE + a.x IN ( + SELECT + /*+ distinct(x) */ x + FROM + $b_sorted + ) +; -- nothing + +SELECT + a.* +FROM + $a AS a +WHERE + a.x IN ( + SELECT + x + FROM + $c_sorted + ) +; -- PruneAdjacentKeys + +SELECT + a.* +FROM + $a AS a +LEFT SEMI JOIN + $b_sorted AS b +ON + a.x == b.x +; -- PruneAdjacentKeys(b_sorted) + +SELECT + a.* +FROM + $b_sorted AS b +RIGHT SEMI JOIN + $a AS a +ON + b.x == a.x +; -- PruneAdjacentKeys(b_sorted) + +SELECT + a.x, + a.t, + b.x +FROM ANY + $a_sorted AS a +JOIN + $b AS b +ON + a.x == b.x +; -- PruneAdjacentKeys(a_sorted) + +SELECT + a.x, + a.t, + b.x +FROM + $a AS a +JOIN ANY + $b_sorted AS b +ON + a.x == b.x +; -- PruneAdjacentKeys(b_sorted) diff --git a/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_select-prune_keys_/formatted.sql b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_select-prune_keys_/formatted.sql index 65a0e7c3a79..80c2fa06403 100644 --- a/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_select-prune_keys_/formatted.sql +++ b/yql/essentials/tests/sql/sql2yql/canondata/test_sql_format.test_select-prune_keys_/formatted.sql @@ -29,6 +29,14 @@ SELECT ListLength(Yql::PruneKeys(AsList(1, 1, 1, 3, 3, 3, 3), $mod2)) ; +SELECT + Yql::PruneAdjacentKeys(AsList(NULL, NULL, NULL, 1, 1, 2, 3, 3, 4, 5), $id) +; + +SELECT + Yql::PruneKeys(AsList(1, NULL, 1, NULL, 1, NULL, 1), $id) +; + -- optimize tests $get_a = ($x) -> { RETURN <|a: $x.a|>; diff --git a/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where.sql b/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where.sql new file mode 100644 index 00000000000..8f23c1f49cd --- /dev/null +++ b/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where.sql @@ -0,0 +1,4 @@ +pragma GroupByExprAfterWhere; +select x from (select 1 as x) +where x = 1 +group by -x as x diff --git a/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where_ver.cfg b/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where_ver.cfg new file mode 100644 index 00000000000..367bc6a9ec0 --- /dev/null +++ b/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where_ver.cfg @@ -0,0 +1 @@ +langver 2025.02 diff --git a/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where_ver.sql b/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where_ver.sql new file mode 100644 index 00000000000..e0542689369 --- /dev/null +++ b/yql/essentials/tests/sql/suites/aggregate/group_by_expr_after_where_ver.sql @@ -0,0 +1,3 @@ +select x from (select 1 as x) +where x = 1 +group by -x as x diff --git a/yql/essentials/tests/sql/suites/join/prune_keys.sql b/yql/essentials/tests/sql/suites/join/prune_keys.sql new file mode 100644 index 00000000000..9f10de9ace7 --- /dev/null +++ b/yql/essentials/tests/sql/suites/join/prune_keys.sql @@ -0,0 +1,50 @@ +pragma config.flags('OptimizerFlags', 'EmitPruneKeys'); + +$a = select * from as_table([ + <|x:1, t:1|>, + <|x:1, t:1|>, + <|x:1, t:2|>, + <|x:3, t:1|>, + <|x:3, t:4|>, + <|x:3, t:2|>, + ]); + +$b = select * from as_table([ + <|x:1, y:1|>, + <|x:1, y:2|>, + <|x:1, y:3|>, + <|x:1, y:3|>, + <|x:2, y:3|>, + <|x:2, y:4|>, + ]); + +$c = select * from as_table([ + <|x:1|>, + <|x:1|>, + <|x:1|>, + <|x:1|>, + <|x:2|>, + <|x:2|>, + ]); + +-- PruneKeys +select a.* from $a as a where a.x in (select x from $b); -- PruneKeys +select a.* from $a as a where a.x in (select /*+ distinct(x) */ x from $b); -- nothing +select a.* from $a as a where a.x in (select x from $c); -- PruneKeys +select a.* from $a as a left semi join $b as b on a.x = b.x; -- PruneKeys(b) +select a.* from $b as b right semi join $a as a on b.x = a.x; -- PruneKeys(b) +select a.x, a.t, b.x from any $a as a join $b as b on a.x == b.x; -- PruneKeys(a) +select a.x, a.t, b.x from $a as a join any $b as b on a.x == b.x; -- PruneKeys(b) + +$a_sorted = select * from $a assume order by x; +$b_sorted = select * from $b assume order by x; +$c_sorted = select * from $c assume order by x; + +-- PruneAdjacentKeys +select a.* from $a as a where a.x in (select x from $b_sorted); -- PruneAdjacentKeys +select a.* from $a as a where a.x in (select /*+ distinct(x) */ x from $b_sorted); -- nothing +select a.* from $a as a where a.x in (select x from $c_sorted); -- PruneAdjacentKeys +select a.* from $a as a left semi join $b_sorted as b on a.x = b.x; -- PruneAdjacentKeys(b_sorted) +select a.* from $b_sorted as b right semi join $a as a on b.x = a.x; -- PruneAdjacentKeys(b_sorted) +select a.x, a.t, b.x from any $a_sorted as a join $b as b on a.x == b.x; -- PruneAdjacentKeys(a_sorted) +select a.x, a.t, b.x from $a as a join any $b_sorted as b on a.x == b.x; -- PruneAdjacentKeys(b_sorted) diff --git a/yql/essentials/tests/sql/suites/select/prune_keys.sql b/yql/essentials/tests/sql/suites/select/prune_keys.sql index 9cf5cb3ced7..f541bcfc61d 100644 --- a/yql/essentials/tests/sql/suites/select/prune_keys.sql +++ b/yql/essentials/tests/sql/suites/select/prune_keys.sql @@ -11,6 +11,9 @@ SELECT Yql::PruneKeys([], $id); $mod2 = ($x) -> { RETURN $x % 2; }; SELECT ListLength(Yql::PruneKeys(AsList(1,1,1,3,3,3,3), $mod2)); +SELECT Yql::PruneAdjacentKeys(AsList(null,null,null,1,1,2,3,3,4,5), $id); +SELECT Yql::PruneKeys(AsList(1,null,1,null,1,null,1), $id); + -- optimize tests $get_a = ($x) -> { RETURN <|a:$x.a|>; }; diff --git a/yt/cpp/mapreduce/http_client/raw_requests.cpp b/yt/cpp/mapreduce/http_client/raw_requests.cpp index 47c0ea204dd..5d19ea1fb9a 100644 --- a/yt/cpp/mapreduce/http_client/raw_requests.cpp +++ b/yt/cpp/mapreduce/http_client/raw_requests.cpp @@ -389,9 +389,7 @@ TNode::TListType LookupRows( fluent.Item("timeout").Value(static_cast<i64>(options.Timeout_->MilliSeconds())); }) .Item("keep_missing_rows").Value(options.KeepMissingRows_) - .DoIf(options.Versioned_.Defined(), [&] (TFluentMap fluent) { - fluent.Item("versioned").Value(*options.Versioned_); - }) + .Item("versioned").Value(options.Versioned_) .DoIf(options.Columns_.Defined(), [&] (TFluentMap fluent) { fluent.Item("column_names").Value(*options.Columns_); }) diff --git a/yt/cpp/mapreduce/interface/client_method_options.h b/yt/cpp/mapreduce/interface/client_method_options.h index 4bb2df112c3..d43020a9e13 100644 --- a/yt/cpp/mapreduce/interface/client_method_options.h +++ b/yt/cpp/mapreduce/interface/client_method_options.h @@ -1019,7 +1019,7 @@ struct TLookupRowsOptions FLUENT_FIELD_DEFAULT(bool, KeepMissingRows, false); /// If set to true returned values will have "timestamp" attribute. - FLUENT_FIELD_OPTION(bool, Versioned); + FLUENT_FIELD_DEFAULT(bool, Versioned, false); }; /// diff --git a/yt/yql/providers/yt/common/yql_configuration.h b/yt/yql/providers/yt/common/yql_configuration.h index 844c89560cc..083504924e0 100644 --- a/yt/yql/providers/yt/common/yql_configuration.h +++ b/yt/yql/providers/yt/common/yql_configuration.h @@ -133,4 +133,5 @@ constexpr bool DEFAULT_ALLOW_REMOTE_CLUSTER_INPUT = false; constexpr bool DEFAULT_USE_COLUMN_GROUPS_FROM_INPUT_TABLE = false; +constexpr bool DEFAULT_USE_NATIVE_DYNAMIC_TABLE_READ = false; } // NYql diff --git a/yt/yql/providers/yt/common/yql_yt_settings.cpp b/yt/yql/providers/yt/common/yql_yt_settings.cpp index 27d6a032a1c..c6d0cb8dec9 100644 --- a/yt/yql/providers/yt/common/yql_yt_settings.cpp +++ b/yt/yql/providers/yt/common/yql_yt_settings.cpp @@ -556,6 +556,7 @@ TYtConfiguration::TYtConfiguration(TTypeAnnotationContext& typeCtx) }); REGISTER_SETTING(*this, _AllowRemoteClusterInput); REGISTER_SETTING(*this, UseColumnGroupsFromInputTables); + REGISTER_SETTING(*this, UseNativeDynamicTableRead); } EReleaseTempDataMode GetReleaseTempDataMode(const TYtSettings& settings) { diff --git a/yt/yql/providers/yt/common/yql_yt_settings.h b/yt/yql/providers/yt/common/yql_yt_settings.h index 5a291b5dab2..677472ecb5d 100644 --- a/yt/yql/providers/yt/common/yql_yt_settings.h +++ b/yt/yql/providers/yt/common/yql_yt_settings.h @@ -317,6 +317,7 @@ struct TYtSettings { NCommon::TConfSetting<bool, false> DropUnusedKeysFromKeyFilter; NCommon::TConfSetting<bool, false> ReportEquiJoinStats; NCommon::TConfSetting<bool, false> UseColumnGroupsFromInputTables; + NCommon::TConfSetting<bool, false> UseNativeDynamicTableRead; }; EReleaseTempDataMode GetReleaseTempDataMode(const TYtSettings& settings); diff --git a/yt/yql/providers/yt/comp_nodes/dq/dq_yt_writer.cpp b/yt/yql/providers/yt/comp_nodes/dq/dq_yt_writer.cpp index 1860f9e3c8f..95436037111 100644 --- a/yt/yql/providers/yt/comp_nodes/dq/dq_yt_writer.cpp +++ b/yt/yql/providers/yt/comp_nodes/dq/dq_yt_writer.cpp @@ -87,7 +87,7 @@ public: NUdf::TUnboxedValuePod DoCalculate(NUdf::TUnboxedValue& state, TComputationContext& ctx) const { if (state.IsFinish()) { - return NUdf::TUnboxedValuePod::MakeFinish(); + return state; } else if (state.IsInvalid()) MakeState(ctx, state); @@ -100,7 +100,7 @@ public: case EFetchResult::Finish: ptr->Finish(); state = NUdf::TUnboxedValuePod::MakeFinish(); - return NUdf::TUnboxedValuePod::MakeFinish(); + return state; } } #ifndef MKQL_DISABLE_CODEGEN diff --git a/yt/yql/providers/yt/gateway/native/yql_yt_native.cpp b/yt/yql/providers/yt/gateway/native/yql_yt_native.cpp index 22e20173d6e..2072158f4e5 100644 --- a/yt/yql/providers/yt/gateway/native/yql_yt_native.cpp +++ b/yt/yql/providers/yt/gateway/native/yql_yt_native.cpp @@ -3385,20 +3385,22 @@ private: } else if (auto limiter = TTableLimiter(range)) { auto entry = execCtx->GetEntry(); bool stop = false; + const bool useNativeDyntableRead = execCtx->Options_.Config()->UseNativeDynamicTableRead.Get().GetOrElse(DEFAULT_USE_NATIVE_DYNAMIC_TABLE_READ); for (size_t i = 0; i < execCtx->InputTables_.size(); ++i) { TString srcTableName = execCtx->InputTables_[i].Name; NYT::TRichYPath srcTable = execCtx->InputTables_[i].Path; - bool isDynamic = execCtx->InputTables_[i].Dynamic; - ui64 recordsCount = execCtx->InputTables_[i].Records; - if (!isDynamic) { - if (!limiter.NextTable(recordsCount)) { - continue; + const bool isDynamic = execCtx->InputTables_[i].Dynamic; + if (!isDynamic || useNativeDyntableRead) { + if (const auto recordsCount = execCtx->InputTables_[i].Records; recordsCount || !isDynamic) { + if (!limiter.NextTable(recordsCount)) { + continue; + } } } else { limiter.NextDynamicTable(); } - if (isDynamic) { + if (isDynamic && !useNativeDyntableRead) { YQL_ENSURE(srcTable.GetRanges().Empty()); stop = NYql::SelectRows(entry->Client, srcTableName, i, specsCache, pullData, limiter); } else { diff --git a/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp b/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp index 8d6b31a1571..c949b461a7b 100644 --- a/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp +++ b/yt/yql/providers/yt/lib/expr_traits/yql_expr_traits.cpp @@ -49,8 +49,11 @@ namespace NYql { (*memoryUsage)["CommonJoinCore"] += FromString<ui64>(memLimitSetting->Child(1)->Content()); } } else if (node.IsCallable("WideCombiner")) { - (*memoryUsage)["WideCombiner"] += FromString<ui64>(node.Child(1U)->Content()); - } else if (NNodes::TCoCombineCore::Match(&node)) { + i64 memLimit = 0LL; + if (TryFromString<i64>(node.Child(1U)->Content(), memLimit)) { + (*memoryUsage)["WideCombiner"] += memLimit; + } + } else if (NNodes::TCoCombineCore::Match(&node) && NNodes::TCoCombineCore::idx_MemLimit < node.ChildrenSize()) { (*memoryUsage)["CombineCore"] += FromString<ui64>(node.Child(NNodes::TCoCombineCore::idx_MemLimit)->Content()); } } diff --git a/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.cpp b/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.cpp index bc785946bd0..bb5c2f4e36a 100644 --- a/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.cpp +++ b/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.cpp @@ -31,6 +31,8 @@ TYtPhysicalOptProposalTransformer::TYtPhysicalOptProposalTransformer(TYtState::T AddHandler(0, &TCoTopSort::Match, HNDL(Sort<true>)); AddHandler(0, &TCoTop::Match, HNDL(Sort<true>)); AddHandler(0, &TYtSort::Match, HNDL(YtSortOverAlreadySorted)); + AddHandler(0, &TCoPruneKeys::Match, HNDL(PushPruneKeysIntoYtOperation)); + AddHandler(0, &TCoPruneAdjacentKeys::Match, HNDL(PushPruneKeysIntoYtOperation)); AddHandler(0, &TCoPartitionByKeyBase::Match, HNDL(PartitionByKey)); AddHandler(0, &TCoFlatMapBase::Match, HNDL(FlatMap)); AddHandler(0, &TCoCombineByKey::Match, HNDL(CombineByKey)); diff --git a/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.h b/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.h index ef783c31c94..cb940e59963 100644 --- a/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.h +++ b/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt.h @@ -149,6 +149,8 @@ private: NNodes::TMaybeNode<NNodes::TExprBase> UpdateDataSourceCluster(NNodes::TExprBase node, TExprContext& ctx) const; + NNodes::TMaybeNode<NNodes::TExprBase> PushPruneKeysIntoYtOperation(NNodes::TExprBase node, TExprContext& ctx) const; + template <typename TLMapType> NNodes::TMaybeNode<NNodes::TExprBase> LMap(NNodes::TExprBase node, TExprContext& ctx) const; diff --git a/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt_misc.cpp b/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt_misc.cpp index 3b65cedd9f9..87b2f66e546 100644 --- a/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt_misc.cpp +++ b/yt/yql/providers/yt/provider/phy_opt/yql_yt_phy_opt_misc.cpp @@ -964,4 +964,74 @@ TMaybeNode<TExprBase> TYtPhysicalOptProposalTransformer::UpdateDataSourceCluster return ctx.ChangeChild(node.Ref(), TYtReadTable::idx_DataSource, MakeDataSource(op.DataSource().Pos(), cluster, ctx).Ptr()); } +TMaybeNode<TExprBase> TYtPhysicalOptProposalTransformer::PushPruneKeysIntoYtOperation(TExprBase node, TExprContext& ctx) const { + auto op = node.Cast<TCoPruneKeysBase>(); + auto extractorLambda = op.Extractor(); + + if (!IsYtProviderInput(op.Input())) { + return node; + } + + TSyncMap syncList; + const ERuntimeClusterSelectionMode selectionMode = + State_->Configuration->RuntimeClusterSelection.Get().GetOrElse(DEFAULT_RUNTIME_CLUSTER_SELECTION); + auto cluster = DeriveClusterFromInput(op.Input(), selectionMode); + if (!cluster || !IsYtCompleteIsolatedLambda(extractorLambda.Ref(), syncList, *cluster, false, selectionMode)) { + return {}; + } + + auto mapper = ctx.Builder(node.Pos()) + .Lambda() + .Param("stream") + .Callable(node.Ref().Content()) + .Arg(0, "stream") + .Add(1, extractorLambda.Ptr()) + .Seal() + .Seal() + .Build(); + + auto outItemType = SilentGetSequenceItemType(op.Input().Ref(), true); + if (!outItemType || !outItemType->IsPersistable()) { + return node; + } + if (!EnsurePersistableYsonTypes(node.Pos(), *outItemType, ctx, State_)) { + return {}; + } + + bool sortedOutput = TCoPruneAdjacentKeys::Match(node.Raw()); + TVector<TYtOutTable> outTables = ConvertOutTablesWithSortAware(mapper, sortedOutput, node.Pos(), + outItemType, ctx, State_, node.Ref().GetConstraintSet()); + + auto settingsBuilder = Build<TCoNameValueTupleList>(ctx, node.Pos()); + if (sortedOutput) { + settingsBuilder + .Add() + .Name() + .Value(ToString(EYtSettingType::Ordered)) + .Build() + .Build(); + } + if (State_->Configuration->UseFlow.Get().GetOrElse(DEFAULT_USE_FLOW)) { + settingsBuilder + .Add() + .Name() + .Value(ToString(EYtSettingType::Flow)) + .Build() + .Build(); + } + + auto map = Build<TYtMap>(ctx, node.Pos()) + .World(GetWorld(op.Input(), {}, ctx)) + .DataSink(MakeDataSink(node.Pos(), *cluster, ctx)) + .Input(ConvertInputTable(op.Input(), ctx)) + .Output() + .Add(outTables) + .Build() + .Settings(settingsBuilder.Done()) + .Mapper(std::move(mapper)) + .Done(); + + return WrapOp(map, ctx); +} + } // namespace NYql diff --git a/yt/yql/tests/sql/suites/join/prune_keys.cfg b/yt/yql/tests/sql/suites/join/prune_keys.cfg new file mode 100644 index 00000000000..9cd81aaa85a --- /dev/null +++ b/yt/yql/tests/sql/suites/join/prune_keys.cfg @@ -0,0 +1,2 @@ +in a_sorted sorted_by_k1.txt +in b_sorted sorted_by_k2.txt diff --git a/yt/yql/tests/sql/suites/join/prune_keys.sql b/yt/yql/tests/sql/suites/join/prune_keys.sql new file mode 100644 index 00000000000..bce76fcb3b4 --- /dev/null +++ b/yt/yql/tests/sql/suites/join/prune_keys.sql @@ -0,0 +1,15 @@ +/* postgres can not */ +use plato; + +pragma yt.JoinMergeTablesLimit = "10"; +pragma config.flags('OptimizerFlags', 'EmitPruneKeys'); + +-- PruneKeys +select * +from a_sorted +where v1 in (select v2 from b_sorted); + +-- PruneAdjacentKeys +select * +from a_sorted +where k1 in (select k2 from b_sorted); diff --git a/yt/yt/client/table_client/schema_serialization_helpers.h b/yt/yt/client/table_client/schema_serialization_helpers.h index f573ac4f4bb..4aeca670370 100644 --- a/yt/yt/client/table_client/schema_serialization_helpers.h +++ b/yt/yt/client/table_client/schema_serialization_helpers.h @@ -4,11 +4,13 @@ #include "schema.h" #include <yt/yt/core/yson/pull_parser.h> + #include <yt/yt/core/ytree/yson_struct.h> namespace NYT::NTableClient { -struct TMaybeDeletedColumnSchema : public TColumnSchema +struct TMaybeDeletedColumnSchema + : public TColumnSchema { DEFINE_BYREF_RO_PROPERTY(std::optional<bool>, Deleted); diff --git a/yt/yt/client/transaction_client/public.h b/yt/yt/client/transaction_client/public.h index 7204cf5808c..033336d6478 100644 --- a/yt/yt/client/transaction_client/public.h +++ b/yt/yt/client/transaction_client/public.h @@ -53,6 +53,7 @@ YT_DEFINE_ERROR_ENUM( ((UnknownClockClusterTag) (11014)) ((ClockClusterTagMismatch) (11015)) ((ChaosCoordinatorsAreNotAvailable) (11016)) + ((NeedLockDynamicTablesBeforeCommit)(11017)) ); //////////////////////////////////////////////////////////////////////////////// diff --git a/yt/yt/core/logging/log_manager.cpp b/yt/yt/core/logging/log_manager.cpp index 6ae91e95cfc..8697b0cd19d 100644 --- a/yt/yt/core/logging/log_manager.cpp +++ b/yt/yt/core/logging/log_manager.cpp @@ -392,6 +392,7 @@ public: if (!IsConfiguredFromEnv()) { DoUpdateConfig(TLogManagerConfig::CreateDefault(), /*fromEnv*/ false); + DefaultConfigured_.store(true); } SystemCategory_ = GetCategory(SystemLoggingCategoryName); @@ -428,6 +429,13 @@ public: if (sync) { future.Get().ThrowOnError(); } + + DefaultConfigured_.store(false); + } + + bool IsDefaultConfigured() + { + return DefaultConfigured_.load(); } void ConfigureFromEnv() @@ -1435,6 +1443,7 @@ private: // Incrementing version forces loggers to update their own default configuration (default level etc.). std::atomic<int> Version_ = 0; + std::atomic<bool> DefaultConfigured_ = false; std::atomic<bool> ConfiguredFromEnv_ = false; // These are just cached (for performance reason) copies from Config_. @@ -1544,6 +1553,14 @@ void TLogManager::Configure(TLogManagerConfigPtr config, bool sync) Impl_->Configure(std::move(config), /*fromEnv*/ false, sync); } +bool TLogManager::IsDefaultConfigured() +{ + [[unlikely]] if (!Impl_->IsInitialized()) { + return false; + } + return Impl_->IsDefaultConfigured(); +} + void TLogManager::ConfigureFromEnv() { [[unlikely]] if (!Impl_->IsInitialized()) { diff --git a/yt/yt/core/logging/log_manager.h b/yt/yt/core/logging/log_manager.h index 37a7073cee2..2f606a124d3 100644 --- a/yt/yt/core/logging/log_manager.h +++ b/yt/yt/core/logging/log_manager.h @@ -32,6 +32,7 @@ public: static TLogManager* Get(); void Configure(TLogManagerConfigPtr config, bool sync = true); + bool IsDefaultConfigured(); void ConfigureFromEnv(); bool IsConfiguredFromEnv(); diff --git a/yt/yt/core/misc/protobuf_helpers-inl.h b/yt/yt/core/misc/protobuf_helpers-inl.h index ed4548cfe40..1d3ce8afd5d 100644 --- a/yt/yt/core/misc/protobuf_helpers-inl.h +++ b/yt/yt/core/misc/protobuf_helpers-inl.h @@ -407,11 +407,11 @@ void FromProtoArrayImpl( originalArray->clear(); originalArray->reserve(serializedArray.size()); for (int i = 0; i < serializedArray.size(); ++i) { - originalArray->emplace( - FromProto<TOriginal>(serializedArray.Get(i))); + originalArray->insert(FromProto<TOriginal>(serializedArray.Get(i))); } } +// Does not check for duplicates. template <class TOriginalKey, class TOriginalValue, class TSerializedArray> void FromProtoArrayImpl( THashMap<TOriginalKey, TOriginalValue>* originalArray, @@ -420,8 +420,7 @@ void FromProtoArrayImpl( originalArray->clear(); originalArray->reserve(serializedArray.size()); for (int i = 0; i < serializedArray.size(); ++i) { - originalArray->emplace( - FromProto<std::pair<TOriginalKey, TOriginalValue>>(serializedArray.Get(i))); + originalArray->insert(FromProto<std::pair<TOriginalKey, TOriginalValue>>(serializedArray.Get(i))); } } @@ -480,6 +479,17 @@ void ToProto( NYT::NDetail::ToProtoArrayImpl(serializedArray, originalArray); } +template <class TKey, class TValue, class TSerializedKey, class TSerializedValue> +void ToProto( + ::google::protobuf::Map<TSerializedKey, TSerializedValue>* serializedMap, + const THashMap<TKey, TValue>& originalMap) +{ + serializedMap->clear(); + for (const auto& [key, value] : originalMap) { + serializedMap->insert(std::pair(ToProto<TSerializedKey>(key), ToProto<TSerializedValue>(value))); + } +} + template <class TOriginalArray, class TSerialized, class... TArgs> void FromProto( TOriginalArray* originalArray, @@ -516,6 +526,19 @@ void CheckedHashSetFromProto( //////////////////////////////////////////////////////////////////////////////// +template <class TKey, class TValue, class TSerializedKey, class TSerializedValue> +void FromProto( + THashMap<TKey, TValue>* originalMap, + const ::google::protobuf::Map<TSerializedKey, TSerializedValue>& serializedMap) +{ + originalMap->clear(); + for (const auto& [serializedKey, serializedValue] : serializedMap) { + EmplaceOrCrash(*originalMap, FromProto<TKey>(serializedKey), FromProto<TValue>(serializedValue)); + } +} + +//////////////////////////////////////////////////////////////////////////////// + template <class TSerialized, NYT::NDetail::CToProtoOriginal TOriginal, class... TArgs> auto ToProto(const TOriginal& original, TArgs&&... args) { diff --git a/yt/yt/core/misc/protobuf_helpers.h b/yt/yt/core/misc/protobuf_helpers.h index c25d7c29023..eff6f223b2f 100644 --- a/yt/yt/core/misc/protobuf_helpers.h +++ b/yt/yt/core/misc/protobuf_helpers.h @@ -123,6 +123,11 @@ void ToProto( ::google::protobuf::RepeatedField<TSerialized>* serializedArray, const TOriginalArray& originalArray); +template <class TKey, class TValue, class TSerializedKey, class TSerializedValue> +void ToProto( + ::google::protobuf::Map<TSerializedKey, TSerializedValue>* serializedMap, + const THashMap<TKey, TValue>& originalMap); + template <class TOriginalArray, class TSerialized, class... TArgs> void FromProto( TOriginalArray* originalArray, @@ -144,6 +149,11 @@ void CheckedHashSetFromProto( THashSet<TOriginal>* originalHashSet, const ::google::protobuf::RepeatedField<TSerialized>& serializedHashSet); +template <class TKey, class TValue, class TSerializedKey, class TSerializedValue> +void FromProto( + THashMap<TKey, TValue>* originalMap, + const ::google::protobuf::Map<TSerializedKey, TSerializedValue>& serializedMap); + //////////////////////////////////////////////////////////////////////////////// template <class TSerialized, class T, class TTag> diff --git a/yt/yt/core/rpc/retrying_channel.cpp b/yt/yt/core/rpc/retrying_channel.cpp index 85a56da1c6c..4d89a3c8ea2 100644 --- a/yt/yt/core/rpc/retrying_channel.cpp +++ b/yt/yt/core/rpc/retrying_channel.cpp @@ -227,7 +227,7 @@ private: })); if (!RetryChecker_.Run(error)) { - ResponseHandler_->HandleError(std::move(error)); + ReportError(error); return; } |
