diff --git a/cpp/test/tools/command_e2e_test.cc b/cpp/test/tools/command_e2e_test.cc index fae33226d..ae70ac386 100644 --- a/cpp/test/tools/command_e2e_test.cc +++ b/cpp/test/tools/command_e2e_test.cc @@ -24,8 +24,13 @@ #include #include +#ifdef _WIN32 +#include +#endif + #include "cli/run_cli.h" #include "cli_test_util.h" +#include "format/input_format.h" namespace { @@ -1624,3 +1629,251 @@ TEST(CliE2E, WriteDistinguishesCsvNullAndEmptyString) { std::remove(csv.c_str()); std::remove(out_path.c_str()); } + +TEST(CliE2E, InvalidUtf8IdentifiersAreReplacedOnlyInOutput) { + storage::libtsfile_init(); + const std::string path = + tsfile_cli_test::unique_temp_path("tsfile_cli_utf8_names", ".tsfile"); + std::string device = "root.device\xff"; + std::string valid_device = "root.\xe4\xb8\xad\xe6\x96\x87"; + const std::string measurement = "value\xe1\x80"; + const std::string valid_measurement = "\xe6\xb8\xa9\xe5\xba\xa6"; + const std::string replacement = "\xef\xbf\xbd"; + { + storage::WriteFile file; + int flags = O_WRONLY | O_CREAT | O_TRUNC; +#ifdef _WIN32 + flags |= O_BINARY; +#endif + ASSERT_EQ(file.create(path, flags, 0666), common::E_OK); + storage::MeasurementSchema schema(measurement, common::INT64, + common::PLAIN, common::UNCOMPRESSED); + storage::MeasurementSchema valid_schema(valid_measurement, + common::INT64, common::PLAIN, + common::UNCOMPRESSED); + storage::TsFileTreeWriter writer(&file); + ASSERT_EQ(writer.register_timeseries(device, &schema), common::E_OK); + ASSERT_EQ(writer.register_timeseries(valid_device, &valid_schema), + common::E_OK); + storage::TsRecord record(device, 0); + record.add_point(measurement, static_cast(42)); + ASSERT_EQ(writer.write(record), common::E_OK); + storage::TsRecord valid_record(valid_device, 0); + valid_record.add_point(valid_measurement, static_cast(7)); + ASSERT_EQ(writer.write(valid_record), common::E_OK); + ASSERT_EQ(writer.flush(), common::E_OK); + ASSERT_EQ(writer.close(), common::E_OK); + } + + const std::string formats[] = {"csv", "ndjson", "table"}; + for (const std::string& format : formats) { + for (const std::string& command : {"ls", "schema", "stats", "count"}) { + SCOPED_TRACE(command + " " + format); + std::ostringstream out, err; + ASSERT_EQ( + tsfile_cli::run_cli({command, "-f", format, path}, out, err), 0) + << err.str(); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(out.str())); + EXPECT_NE(out.str().find("root.device" + replacement), + std::string::npos); + EXPECT_NE(out.str().find(valid_device), std::string::npos); + } + for (const std::string& command : {"head", "cat"}) { + SCOPED_TRACE(command + " " + format); + std::ostringstream out, err; + ASSERT_EQ(tsfile_cli::run_cli({command, "-d", device, "-m", + measurement, "-f", format, path}, + out, err), + 0) + << err.str(); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(out.str())); + EXPECT_NE(out.str().find("value" + replacement), std::string::npos); + EXPECT_NE(out.str().find("42"), std::string::npos); + } + + const std::string output = tsfile_cli_test::unique_temp_path( + "tsfile_cli_utf8_export", "." + format); + std::ostringstream out, err; + ASSERT_EQ(tsfile_cli::run_cli({"export", "-d", device, "--type", format, + "-o", output, path}, + out, err), + 0) + << err.str(); + std::ifstream file(output.c_str(), std::ios::binary); + std::ostringstream content; + content << file.rdbuf(); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(content.str())); + EXPECT_NE(content.str().find("value" + replacement), std::string::npos); + file.close(); + std::remove(output.c_str()); + } + + std::ostringstream sketch, sketch_err; + ASSERT_EQ(tsfile_cli::run_cli({"sketch", path}, sketch, sketch_err), 0) + << sketch_err.str(); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(sketch.str())); + EXPECT_NE(sketch.str().find("value" + replacement), std::string::npos); + + const std::string dir = + tsfile_cli_test::unique_temp_path("tsfile_cli_utf8_manifest", ""); + std::ostringstream out, err; + ASSERT_EQ( + tsfile_cli::run_cli({"export", "-d", device, "-d", valid_device, + "--type", "ndjson", "--output-dir", dir, path}, + out, err), + 0) + << err.str(); + std::ifstream manifest((dir + "/_manifest.json").c_str(), std::ios::binary); + std::ostringstream content; + content << manifest.rdbuf(); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(content.str())); + EXPECT_NE(content.str().find("root.device" + replacement), + std::string::npos); + EXPECT_NE(content.str().find(valid_device), std::string::npos); + manifest.close(); + std::remove((dir + "/0001.ndjson").c_str()); + std::remove((dir + "/0002.ndjson").c_str()); + std::remove((dir + "/_manifest.json").c_str()); +#ifdef _WIN32 + _rmdir(dir.c_str()); +#else + rmdir(dir.c_str()); +#endif + std::remove(path.c_str()); +} + +TEST(CliE2E, Utf8JsonKeyCollisionsFailWithoutPublishingOutput) { + storage::libtsfile_init(); + const std::string path = tsfile_cli_test::unique_temp_path( + "tsfile_cli_utf8_collision", ".tsfile"); + std::string device = "root.collision"; + const std::vector names = {"bad\xff", "bad\xfe", + "bad\xef\xbf\xbd"}; + { + storage::WriteFile file; + int flags = O_WRONLY | O_CREAT | O_TRUNC; +#ifdef _WIN32 + flags |= O_BINARY; +#endif + ASSERT_EQ(file.create(path, flags, 0666), common::E_OK); + storage::MeasurementSchema first(names[0], common::INT64, common::PLAIN, + common::UNCOMPRESSED); + storage::MeasurementSchema second(names[1], common::INT64, + common::PLAIN, common::UNCOMPRESSED); + storage::MeasurementSchema third(names[2], common::INT64, common::PLAIN, + common::UNCOMPRESSED); + storage::TsFileTreeWriter writer(&file); + ASSERT_EQ(writer.register_timeseries(device, &first), common::E_OK); + ASSERT_EQ(writer.register_timeseries(device, &second), common::E_OK); + ASSERT_EQ(writer.register_timeseries(device, &third), common::E_OK); + storage::TsRecord record(device, 0); + record.add_point(names[0], static_cast(42)); + record.add_point(names[1], static_cast(7)); + record.add_point(names[2], static_cast(9)); + ASSERT_EQ(writer.write(record), common::E_OK); + ASSERT_EQ(writer.flush(), common::E_OK); + ASSERT_EQ(writer.close(), common::E_OK); + } + + for (const std::string& command : {"head", "cat"}) { + for (size_t second : {size_t(1), size_t(2)}) { + for (bool empty : {false, true}) { + SCOPED_TRACE(command + " second=" + std::to_string(second) + + " empty=" + std::to_string(empty)); + std::vector args = { + command, "-d", device, "-m", names[0], + "-m", names[second], "-f", "ndjson", path}; + if (empty) { + args.insert(args.end() - 1, {"--start", "0", "-n", "0"}); + } + std::ostringstream out, err; + EXPECT_EQ(tsfile_cli::run_cli(args, out, err), 3) << err.str(); + EXPECT_TRUE(out.str().empty()); + EXPECT_NE(err.str().find("duplicate NDJSON column name"), + std::string::npos); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(err.str())); + } + } + } + + const std::string output = tsfile_cli_test::unique_temp_path( + "tsfile_cli_utf8_collision_export", ".ndjson"); + for (bool existing : {false, true}) { + if (existing) { + std::ofstream target(output.c_str(), std::ios::binary); + target << "keep existing output\n"; + } + std::vector args = {"export", "-d", device, "--type", + "ndjson", "-o", output, path}; + if (existing) { + args.insert(args.end() - 1, "--force"); + } + std::ostringstream out, err; + EXPECT_EQ(tsfile_cli::run_cli(args, out, err), 3) << err.str(); + EXPECT_TRUE(out.str().empty()); + EXPECT_NE(err.str().find("duplicate NDJSON column name"), + std::string::npos); + std::ifstream target(output.c_str(), std::ios::binary); + if (existing) { + ASSERT_TRUE(target.is_open()); + std::ostringstream content; + content << target.rdbuf(); + EXPECT_EQ(content.str(), "keep existing output\n"); + } else { + EXPECT_FALSE(target.is_open()); + } + } + std::remove(output.c_str()); + std::remove(path.c_str()); +} + +TEST(CliE2E, TableUtf8JsonKeyCollisionsFailInRowsAndStats) { + storage::libtsfile_init(); + const std::string path = tsfile_cli_test::unique_temp_path( + "tsfile_cli_utf8_table_collision", ".tsfile"); + std::string table = "collision"; + const std::vector names = {"bad\xff", "bad\xfe", "value"}; + { + storage::WriteFile file; + int flags = O_WRONLY | O_CREAT | O_TRUNC; +#ifdef _WIN32 + flags |= O_BINARY; +#endif + ASSERT_EQ(file.create(path, flags, 0666), common::E_OK); + storage::TableSchema schema( + table, {common::ColumnSchema(names[0], common::STRING, + common::UNCOMPRESSED, common::PLAIN, + common::ColumnCategory::TAG), + common::ColumnSchema(names[1], common::STRING, + common::UNCOMPRESSED, common::PLAIN, + common::ColumnCategory::TAG), + common::ColumnSchema(names[2], common::INT64, + common::UNCOMPRESSED, common::PLAIN, + common::ColumnCategory::FIELD)}); + storage::TsFileTableWriter writer(&file, &schema); + storage::Tablet tablet( + table, names, {common::STRING, common::STRING, common::INT64}, + {common::ColumnCategory::TAG, common::ColumnCategory::TAG, + common::ColumnCategory::FIELD}, + 1); + tablet.add_timestamp(0, static_cast(0)); + tablet.add_value(0, names[0], "first"); + tablet.add_value(0, names[1], "second"); + tablet.add_value(0, names[2], static_cast(42)); + ASSERT_EQ(writer.write_table(tablet), common::E_OK); + ASSERT_EQ(writer.flush(), common::E_OK); + ASSERT_EQ(writer.close(), common::E_OK); + } + for (const std::string& command : {"head", "cat", "stats"}) { + SCOPED_TRACE(command); + std::ostringstream out, err; + EXPECT_EQ(tsfile_cli::run_cli( + {command, "-t", table, "-f", "ndjson", path}, out, err), + 3); + EXPECT_TRUE(out.str().empty()); + EXPECT_NE(err.str().find("duplicate NDJSON column name"), + std::string::npos); + EXPECT_TRUE(tsfile_cli::is_valid_utf8(err.str())); + } + std::remove(path.c_str()); +} diff --git a/cpp/test/tools/output_format_test.cc b/cpp/test/tools/output_format_test.cc index 55324ff3f..b58fafe3d 100644 --- a/cpp/test/tools/output_format_test.cc +++ b/cpp/test/tools/output_format_test.cc @@ -23,6 +23,7 @@ #include #include +#include #include #include "common/db_common.h" @@ -93,6 +94,56 @@ TEST(JsonEscapeTest, EscapesQuotesBackslashAndControls) { EXPECT_EQ(tsfile_cli::json_escape("tab\there"), "tab\\there"); } +TEST(Utf8OutputTest, PreservesValidSequencesAcrossFormats) { + const std::string text = + "ASCII\xc2\x80\xdf\xbf\xe0\xa0\x80\xe4\xb8\xad" + "\xed\x9f\xbf\xee\x80\x80\xef\xbb\xbf\xf0\x90\x80\x80" + "\xf0\x9f\x98\x80\xf4\x8f\xbf\xbf"; + EXPECT_EQ(tsfile_cli::csv_escape(text), text); + EXPECT_EQ(tsfile_cli::json_escape(text), text); + EXPECT_EQ(tsfile_cli::table_escape(text), text); +} + +TEST(Utf8OutputTest, ReplacesMaximalInvalidSubpartsAcrossFormats) { + const std::string replacement = "\xef\xbf\xbd"; + const std::vector> cases = { + {"\xff", replacement}, + {"\x80\xbf", replacement + replacement}, + {"\xc0\xaf", replacement + replacement}, + {"\xe0\x80\xbf", replacement + replacement + replacement}, + {"\xed\xa0\x80", replacement + replacement + replacement}, + {"\xf4\x91\x92\x93", + replacement + replacement + replacement + replacement}, + {"\xc2", replacement}, + {"\xe1\x80", replacement}, + {"\xf0\x91\x92", replacement}, + {"a\xc2" + "b", + "a" + replacement + "b"}, + {"\xe1\x80" + "A", + replacement + "A"}, + {"\xf1\xbf\xe4\xb8\xad", replacement + "\xe4\xb8\xad"}, + {"\xe1\x80\xe2\xf0\x91\x92\xf1\xbf" + "A", + replacement + replacement + replacement + replacement + "A"}, + }; + for (const auto& test : cases) { + SCOPED_TRACE(::testing::PrintToString(test.first)); + EXPECT_EQ(tsfile_cli::csv_escape(test.first), test.second); + EXPECT_EQ(tsfile_cli::json_escape(test.first), test.second); + EXPECT_EQ(tsfile_cli::table_escape(test.first), test.second); + } +} + +TEST(Utf8OutputTest, PreservesDelimitersAfterAnInvalidSequence) { + const std::string text = "\xe1\x80\",\n"; + const std::string replacement = "\xef\xbf\xbd"; + EXPECT_EQ(tsfile_cli::csv_escape(text), "\"" + replacement + "\"\",\n\""); + EXPECT_EQ(tsfile_cli::json_escape(text), replacement + "\\\",\\n"); + EXPECT_EQ(tsfile_cli::table_escape(text), replacement + "\",\\n"); +} + TEST(TableEscapeTest, EscapesBackslashAndNamedControls) { EXPECT_EQ(tsfile_cli::table_escape("a\\b"), "a\\\\b"); EXPECT_EQ(tsfile_cli::table_escape("line\nbreak"), "line\\nbreak"); @@ -117,8 +168,7 @@ TEST(TableEscapeTest, ControlByteBoundaries) { // Space (0x20) is printable: it must survive verbatim, never be escaped. EXPECT_EQ(tsfile_cli::table_escape("a b"), "a b"); - // Bytes >= 0x80 are UTF-8 lead/continuation bytes; escaping them would - // corrupt multi-byte sequences, so they must pass through unchanged. + // Well-formed UTF-8 must pass through unchanged. EXPECT_EQ(tsfile_cli::table_escape("\xe4\xb8\xad"), "\xe4\xb8\xad"); } @@ -264,3 +314,59 @@ TEST(RowWriterTest, ReportsFlushFailure) { ASSERT_TRUE(writer.write({"value"}, {false})); EXPECT_FALSE(writer.finish()); } + +TEST(RowWriterTest, ReplacesInvalidUtf8InHeadersAndValues) { + const std::string header = "name\xff"; + const std::string value = "value\xe1\x80"; + const std::string replacement = "\xef\xbf\xbd"; + const OutputFormat formats[] = {OutputFormat::kCsv, OutputFormat::kJson, + OutputFormat::kTable}; + for (OutputFormat format : formats) { + std::ostringstream out; + RowWriter writer(out, format, {header}, {common::STRING}, false); + ASSERT_TRUE(writer.write({value}, {false})); + ASSERT_TRUE(writer.finish()); + const std::string expected = + format == OutputFormat::kJson + ? "{\"name" + replacement + "\":\"value" + replacement + "\"}\n" + : "name" + replacement + "\nvalue" + replacement + "\n"; + EXPECT_EQ(out.str(), expected); + } +} + +TEST(RowWriterTest, PreservesNonUtf8BlobBytesAsHex) { + const OutputFormat formats[] = {OutputFormat::kCsv, OutputFormat::kJson, + OutputFormat::kTable}; + for (OutputFormat format : formats) { + std::ostringstream out; + RowWriter writer(out, format, {"payload"}, {common::BLOB}, false); + ASSERT_TRUE(writer.write({std::string("\xff\0\x80", 3)}, {false})); + ASSERT_TRUE(writer.finish()); + EXPECT_NE(out.str().find("0xff0080"), std::string::npos); + } +} + +TEST(RowWriterTest, RejectsJsonKeyCollisionsAfterUtf8Replacement) { + const std::vector> headers = { + {"bad\xff", "bad\xfe"}, + {"bad\xff", "bad\xef\xbf\xbd"}, + }; + for (const auto& header : headers) { + for (bool no_header : {false, true}) { + std::ostringstream out; + RowWriter writer(out, OutputFormat::kJson, header, + {common::INT64, common::INT64}, no_header); + EXPECT_FALSE(writer.write({"42", "7"}, {false, false})); + EXPECT_FALSE(writer.finish()); + EXPECT_TRUE(out.str().empty()); + } + } +} + +TEST(RowWriterTest, RejectsJsonKeyCollisionsWithoutRows) { + std::ostringstream out; + RowWriter writer(out, OutputFormat::kJson, {"bad\xff", "bad\xfe"}, + {common::STRING, common::STRING}, false); + EXPECT_FALSE(writer.finish()); + EXPECT_TRUE(out.str().empty()); +} diff --git a/cpp/tools/README.md b/cpp/tools/README.md index 37f704896..874f52bb0 100644 --- a/cpp/tools/README.md +++ b/cpp/tools/README.md @@ -121,6 +121,12 @@ values. The `table` format uses a temporary spool to align columns with bounded memory; prefer `csv`/`ndjson` when temporary disk use is undesirable. `sketch` does not accept `--format`. +Text output preserves well-formed UTF-8 and replaces malformed UTF-8 in names +and text with U+FFFD (`�`). This also applies to sketch output and export +manifests. BLOB values remain hexadecimal strings. +NDJSON output fails before writing rows if replacement would produce duplicate +column names. + ```bash BIN=cpp/build/Debug/bin/tsfile-cli $BIN ls -f csv data.tsfile # list tables / devices diff --git a/cpp/tools/commands/cmd_export.cc b/cpp/tools/commands/cmd_export.cc index ecca955e1..9a7a1c7ee 100644 --- a/cpp/tools/commands/cmd_export.cc +++ b/cpp/tools/commands/cmd_export.cc @@ -37,6 +37,7 @@ #include "common/device_id.h" #include "common/schema.h" #include "format/atomic_output.h" +#include "format/output_format.h" #include "reader/tsfile_reader.h" namespace tsfile_cli { @@ -143,32 +144,6 @@ std::string numbered_file_name(size_t index, ParsedArgs::Format fmt) { return std::string(buf) + extension_for_format(fmt); } -std::string json_escape(const std::string& s) { - std::ostringstream out; - for (char c : s) { - switch (c) { - case '\\': - out << "\\\\"; - break; - case '"': - out << "\\\""; - break; - case '\n': - out << "\\n"; - break; - case '\r': - out << "\\r"; - break; - case '\t': - out << "\\t"; - break; - default: - out << c; - } - } - return out.str(); -} - struct ManifestEntry { std::string file; std::string model; diff --git a/cpp/tools/commands/cmd_sketch.cc b/cpp/tools/commands/cmd_sketch.cc index eb762048f..eff543fb2 100644 --- a/cpp/tools/commands/cmd_sketch.cc +++ b/cpp/tools/commands/cmd_sketch.cc @@ -128,8 +128,7 @@ int write_atomic_text(const std::string& path, const std::string& content, return code; } { - std::ofstream output(tmp.c_str(), - std::ios::binary | std::ios::trunc); + std::ofstream output(tmp.c_str(), std::ios::binary | std::ios::trunc); if (!output.is_open()) { err << "Error: cannot create output target '" << path << "'\n"; remove_atomic_temp(tmp, err); @@ -460,8 +459,7 @@ class SketchPrinter { to_string_u32(layout_.footer_.bloom_filter_hash_count)); std::ostringstream bloom; bloom << "[Bloom Filter] , filterCapacity=" - << layout_.footer_.bloom_filter_size - << ", hashFunctionSize=" + << layout_.footer_.bloom_filter_size << ", hashFunctionSize=" << layout_.footer_.bloom_filter_hash_count; print_line(out, layout_.footer_.bloom_filter_offset, bloom.str()); } @@ -587,21 +585,16 @@ int cmd_sketch(const ParsedArgs& args, std::ostream& out, std::ostream& err) { SketchPrinter printer; std::ostringstream content; int code = printer.run(args, content, err); - if (code != kExitOk) { - if (args.output.empty()) { - out << content.str(); - out.flush(); - return out.good() ? code : kExitRuntime; - } + if (code != kExitOk && !args.output.empty()) { return code; } + const std::string text = replace_invalid_utf8(content.str()); if (!args.output.empty()) { - return write_atomic_text(args.output, content.str(), args.file, - args.force, err); + return write_atomic_text(args.output, text, args.file, args.force, err); } - out << content.str(); + out << text; out.flush(); - return out.good() ? kExitOk : kExitRuntime; + return out.good() ? code : kExitRuntime; } } // namespace tsfile_cli diff --git a/cpp/tools/commands/cmd_stats.cc b/cpp/tools/commands/cmd_stats.cc index 01bd5ba57..df5371f1b 100644 --- a/cpp/tools/commands/cmd_stats.cc +++ b/cpp/tools/commands/cmd_stats.cc @@ -500,6 +500,10 @@ int cmd_table_stats(const ParsedArgs& args, storage::TsFileReader& reader, types.push_back(common::STRING); } RowWriter w(out, fmt, headers, types, args.no_header); + if (!w.error().empty()) { + err << "Error: " << w.error() << "\n"; + return kExitRuntime; + } for (const TableStatsSummary& summary : summaries) { std::map local_tag_positions; for (size_t i = 0; i < summary.tag_indexes.size(); ++i) { diff --git a/cpp/tools/commands/row_query.cc b/cpp/tools/commands/row_query.cc index 828bb6c55..f6150600f 100644 --- a/cpp/tools/commands/row_query.cc +++ b/cpp/tools/commands/row_query.cc @@ -361,11 +361,17 @@ int run_row_query(const ParsedArgs& args, storage::TsFileReader& reader, // that callers could mistake for a complete result. The final write is // still checked separately so stdout errors remain runtime failures. std::ostringstream staged; - int wret = push_down ? emit_result_set(rs, fmt, args.no_header, staged, 0, - -1, emitted_rows) - : emit_result_set(rs, fmt, args.no_header, staged, - offset, limit, emitted_rows); + std::string output_error; + int wret = push_down + ? emit_result_set(rs, fmt, args.no_header, staged, 0, -1, + emitted_rows, &output_error) + : emit_result_set(rs, fmt, args.no_header, staged, offset, + limit, emitted_rows, &output_error); reader.destroy_query_data_set(rs); + if (!output_error.empty()) { + err << "Error: " << output_error << "\n"; + return kExitRuntime; + } if (wret == common::E_OK) { const std::string bytes = staged.str(); out.write(bytes.data(), static_cast(bytes.size())); diff --git a/cpp/tools/format/output_format.cc b/cpp/tools/format/output_format.cc index 63aa39ad6..b6ba39c1b 100644 --- a/cpp/tools/format/output_format.cc +++ b/cpp/tools/format/output_format.cc @@ -21,6 +21,7 @@ #include #include +#include #include #include "common/csv_utils.h" @@ -192,14 +193,72 @@ const char* compression_name(common::CompressionType c) { } } +namespace { + +// UTF-8 encoding of the replacement character U+FFFD. +constexpr unsigned char kUtf8Replacement[] = {0xef, 0xbf, 0xbd}; + +template +void for_each_utf8_byte(const std::string& s, Emit emit) { + size_t i = 0; + while (i < s.size()) { + const unsigned char lead = static_cast(s[i]); + size_t length = 0; + if (lead <= 0x7f) { + length = 1; + } else if (lead >= 0xc2 && lead <= 0xdf) { + length = 2; + } else if (lead >= 0xe0 && lead <= 0xef) { + length = 3; + } else if (lead >= 0xf0 && lead <= 0xf4) { + length = 4; + } + size_t consumed = 1; + while (consumed < length && consumed < s.size() - i) { + const unsigned char next = + static_cast(s[i + consumed]); + if (next < 0x80 || next > 0xbf || + (consumed == 1 && ((lead == 0xe0 && next < 0xa0) || + (lead == 0xed && next > 0x9f) || + (lead == 0xf0 && next < 0x90) || + (lead == 0xf4 && next > 0x8f)))) { + break; + } + ++consumed; + } + if (consumed == length) { + for (size_t j = 0; j < consumed; ++j) { + emit(static_cast(s[i + j])); + } + } else { + for (unsigned char c : kUtf8Replacement) { + emit(c); + } + } + // Consume only a valid prefix of this character, leaving subsequent + // ASCII or valid UTF-8 for the next iteration, even after truncation. + i += consumed; + } +} + +} // namespace + +std::string replace_invalid_utf8(const std::string& s) { + std::string out; + out.reserve(s.size()); + for_each_utf8_byte( + s, [&out](unsigned char c) { out += static_cast(c); }); + return out; +} + std::string csv_escape(const std::string& field) { - return common::csv_escape(field, ','); + return common::csv_escape(replace_invalid_utf8(field), ','); } std::string json_escape(const std::string& s) { std::string out; out.reserve(s.size() + 2); - for (unsigned char c : s) { + for_each_utf8_byte(s, [&out](unsigned char c) { switch (c) { case '"': out += "\\\""; @@ -231,7 +290,7 @@ std::string json_escape(const std::string& s) { out += static_cast(c); } } - } + }); return out; } @@ -243,7 +302,7 @@ std::string json_escape(const std::string& s) { std::string table_escape(const std::string& s) { std::string out; out.reserve(s.size() + 2); - for (unsigned char c : s) { + for_each_utf8_byte(s, [&out](unsigned char c) { switch (c) { case '\\': out += "\\\\"; @@ -260,8 +319,8 @@ std::string table_escape(const std::string& s) { default: // C0 controls (0x00-0x1f) plus DEL (0x7f) are the complete // set of non-printable control bytes per C iscntrl(). Bytes - // >= 0x80 are UTF-8 continuation/lead bytes and must pass - // through untouched so multi-byte sequences survive. + // >= 0x80 now belong to well-formed UTF-8 and pass through + // untouched so multi-byte sequences survive. if (c < 0x20 || c == 0x7f) { char buf[8]; std::snprintf(buf, sizeof(buf), "\\u%04x", c); @@ -270,7 +329,7 @@ std::string table_escape(const std::string& s) { out += static_cast(c); } } - } + }); return out; } @@ -308,7 +367,21 @@ RowWriter::RowWriter(std::ostream& out, OutputFormat fmt, types_(std::move(types)), no_header_(no_header), table_widths_(header_.size(), 0) { - if (!no_header_) { + if (fmt_ == OutputFormat::kJson) { + std::set keys; + json_keys_.reserve(header_.size()); + for (const std::string& name : header_) { + std::string key = json_escape(name); + if (!keys.insert(key).second) { + error_ = + "duplicate NDJSON column name after UTF-8 replacement: \"" + + key + "\""; + return; + } + json_keys_.push_back(std::move(key)); + } + } + if (fmt_ == OutputFormat::kTable && !no_header_) { for (size_t i = 0; i < header_.size(); ++i) { table_widths_[i] = table_escape(header_[i]).size(); } @@ -372,7 +445,7 @@ bool RowWriter::write(const std::vector& cells, bool RowWriter::write(const std::vector& cells, const std::vector& is_null, const std::vector& row_types) { - if (!out_.good()) { + if (!error_.empty() || !out_.good()) { return false; } if (fmt_ == OutputFormat::kTable) { @@ -410,7 +483,7 @@ bool RowWriter::write(const std::vector& cells, if (i) { out_ << ","; } - out_ << "\"" << json_escape(header_[i]) << "\":"; + out_ << "\"" << json_keys_[i] << "\":"; const common::TSDataType type = i < row_types.size() ? row_types[i] : common::STRING; if (i < is_null.size() && is_null[i]) { @@ -463,7 +536,7 @@ bool RowWriter::write(const std::vector& cells, } bool RowWriter::finish() { - if (!out_.good()) { + if (!error_.empty() || !out_.good()) { return false; } if (fmt_ != OutputFormat::kTable) { diff --git a/cpp/tools/format/output_format.h b/cpp/tools/format/output_format.h index 1a9ed52e9..02d005b41 100644 --- a/cpp/tools/format/output_format.h +++ b/cpp/tools/format/output_format.h @@ -42,6 +42,8 @@ const char* tsdatatype_name(common::TSDataType t); const char* tsencoding_name(common::TSEncoding e); const char* compression_name(common::CompressionType c); +// Replace each maximal ill-formed UTF-8 subpart with U+FFFD (Unicode 3.9.6). +std::string replace_invalid_utf8(const std::string& s); std::string csv_escape(const std::string& field); std::string json_escape(const std::string& s); std::string table_escape(const std::string& s); @@ -60,6 +62,9 @@ class RowWriter { const std::vector& row_types); bool finish(); + // Header validation failures are reported before any output is written. + const std::string& error() const { return error_; } + private: bool ensure_header(); bool emits_json_bare(common::TSDataType type) const; @@ -69,6 +74,8 @@ class RowWriter { std::ostream& out_; OutputFormat fmt_; std::vector header_; + std::vector json_keys_; + std::string error_; std::vector types_; bool no_header_; bool header_done_ = false; diff --git a/cpp/tools/format/result_set_format.cc b/cpp/tools/format/result_set_format.cc index ea0b350b4..a67e39c8c 100644 --- a/cpp/tools/format/result_set_format.cc +++ b/cpp/tools/format/result_set_format.cc @@ -74,7 +74,10 @@ std::string cell_to_string(storage::ResultSet* rs, uint32_t i, int emit_result_set(storage::ResultSet* rs, OutputFormat fmt, bool no_header, std::ostream& out, long long offset, long long limit, - long long* emitted_rows) { + long long* emitted_rows, std::string* output_error) { + if (output_error != nullptr) { + output_error->clear(); + } auto meta = rs->get_metadata(); const uint32_t ncol = meta->get_column_count(); std::vector header; @@ -87,6 +90,12 @@ int emit_result_set(storage::ResultSet* rs, OutputFormat fmt, bool no_header, } RowWriter writer(out, fmt, header, types, no_header); + if (!writer.error().empty()) { + if (output_error != nullptr) { + *output_error = writer.error(); + } + return common::E_INVALID_ARG; + } bool has_next = false; int code = common::E_OK; long long skipped = 0; diff --git a/cpp/tools/format/result_set_format.h b/cpp/tools/format/result_set_format.h index 40814a47b..84ea8f0e9 100644 --- a/cpp/tools/format/result_set_format.h +++ b/cpp/tools/format/result_set_format.h @@ -32,9 +32,11 @@ namespace tsfile_cli { std::string cell_to_string(storage::ResultSet* rs, uint32_t col_index, common::TSDataType type); +// output_error receives a diagnostic when output column names are invalid. int emit_result_set(storage::ResultSet* rs, OutputFormat fmt, bool no_header, std::ostream& out, long long offset = 0, - long long limit = -1, long long* emitted_rows = nullptr); + long long limit = -1, long long* emitted_rows = nullptr, + std::string* output_error = nullptr); int emit_result_set_sampled(storage::ResultSet* rs, OutputFormat fmt, bool no_header, std::ostream& out, long long limit,