diff --git a/elasticgraph-indexer/lib/elastic_graph/indexer.rb b/elasticgraph-indexer/lib/elastic_graph/indexer.rb index 218b6818d..3bb1a34e6 100644 --- a/elasticgraph-indexer/lib/elastic_graph/indexer.rb +++ b/elasticgraph-indexer/lib/elastic_graph/indexer.rb @@ -85,6 +85,7 @@ def operation_factory record_preparer_factory: record_preparer_factory, logger: datastore_core.logger, skip_derived_indexing_type_updates: config.skip_derived_indexing_type_updates, + skip_record_validation_percents_by_type: config.skip_record_validation_percents_by_type, configure_record_validator: nil ) end diff --git a/elasticgraph-indexer/lib/elastic_graph/indexer/config.rb b/elasticgraph-indexer/lib/elastic_graph/indexer/config.rb index 6eca6a902..4af8d1bf2 100644 --- a/elasticgraph-indexer/lib/elastic_graph/indexer/config.rb +++ b/elasticgraph-indexer/lib/elastic_graph/indexer/config.rb @@ -11,7 +11,7 @@ module ElasticGraph class Indexer - class Config < Support::Config.define(:latency_slo_thresholds_by_timestamp_in_ms, :skip_derived_indexing_type_updates, :extension_modules) + class Config < Support::Config.define(:latency_slo_thresholds_by_timestamp_in_ms, :skip_derived_indexing_type_updates, :skip_record_validation_percents_by_type, :extension_modules) json_schema at: "indexer", optional: false, description: "Configuration for indexing operations and metrics used by `elasticgraph-indexer`.", @@ -43,15 +43,45 @@ class Config < Support::Config.define(:latency_slo_thresholds_by_timestamp_in_ms {"WidgetWorkspace" => ["ABC12345678"]} ] }, + skip_record_validation_percents_by_type: { + description: "Map of GraphQL type names to the percentage of records of that type whose per-record " \ + "JSON schema validation should be skipped. `0` (or an absent key) validates every record of the " \ + "type; `100` skips every record; values in between sample, and may be fractional. The decision is " \ + "deterministic per event id (`type:id@vversion`), so the same event makes the same choice on every " \ + "retry and on every indexer pod. The event envelope (op, id, type, version, json_schema_version, " \ + "latency_timestamps) is always validated, regardless of this setting.\n\n" \ + "With a large schema the per-record schema walk consumes a significant share of indexing CPU: every " \ + "record is checked against every regex, enum, min/max, format, and abstract-type discriminator " \ + "defined for its type. Skipping it trades that check for throughput, which is worthwhile when " \ + "backfilling data that was already validated upstream. Leaving a percentage of records validated " \ + "keeps a canary in place so schema drift still surfaces.\n\n" \ + "Note: skipping validation makes malformed-data detection later and less precise. A malformation " \ + "found while building an event's operations is still reported as an isolated event failure, carrying " \ + "the message validation itself would have produced. But one found only while serializing an " \ + "operation for the datastore (an unparsable rollover index timestamp, or a missing custom routing " \ + "field) raises an error that fails the entire batch, including the well-formed events in it. Since " \ + "such a batch produces no partial-failure response, the queue redelivers all of its events, and the " \ + "malformed record fails them again on each retry until it is drained to the dead letter queue. Leave " \ + "this empty for live-traffic ingestion.", + type: "object", + patternProperties: {/^[A-Z]\w*$/.source => {type: "number", minimum: 0, maximum: 100}}, + additionalProperties: false, + default: {}, # : untyped + examples: [ + {}, # : untyped + {"Widget" => 90, "Component" => 100} + ] + }, extension_modules: Support::Config::EXTENSION_MODULE_SCHEMA } private - def convert_values(skip_derived_indexing_type_updates:, latency_slo_thresholds_by_timestamp_in_ms:, extension_modules:) + def convert_values(skip_derived_indexing_type_updates:, latency_slo_thresholds_by_timestamp_in_ms:, skip_record_validation_percents_by_type:, extension_modules:) { skip_derived_indexing_type_updates: skip_derived_indexing_type_updates.transform_values(&:to_set), latency_slo_thresholds_by_timestamp_in_ms: latency_slo_thresholds_by_timestamp_in_ms, + skip_record_validation_percents_by_type: skip_record_validation_percents_by_type.transform_values(&:to_f), extension_modules: SchemaArtifacts::RuntimeMetadata::ExtensionLoader.load_component_extensions(extension_modules) } end diff --git a/elasticgraph-indexer/lib/elastic_graph/indexer/operation/factory.rb b/elasticgraph-indexer/lib/elastic_graph/indexer/operation/factory.rb index 96c5dff61..1b3d52f00 100644 --- a/elasticgraph-indexer/lib/elastic_graph/indexer/operation/factory.rb +++ b/elasticgraph-indexer/lib/elastic_graph/indexer/operation/factory.rb @@ -13,6 +13,7 @@ require "elastic_graph/indexer/record_preparer" require "elastic_graph/support/json_schema/validator_factory" require "elastic_graph/support/memoizable_data" +require "zlib" module ElasticGraph class Indexer @@ -23,6 +24,7 @@ class Factory < Support::MemoizableData.define( :record_preparer_factory, :logger, :skip_derived_indexing_type_updates, + :skip_record_validation_percents_by_type, :configure_record_validator ) def build(event) @@ -40,15 +42,44 @@ def build(event) return build_failed_result(event, "event payload", error_message) end - failed_result = validate_record_returning_failure(event, selected_json_schema_version) - failed_result || BuildResult.success(build_all_operations_for( - event, - record_preparer_factory.for_json_schema_version(selected_json_schema_version) - )) + graphql_type_name = event.fetch("type") + + if skip_validation?(graphql_type_name, event) + build_success_result_isolating_malformed_records(event, graphql_type_name, selected_json_schema_version) + else + validate_record_returning_failure(event, graphql_type_name, selected_json_schema_version) || + build_success_result(event, selected_json_schema_version, type_with_skipped_validation: nil) + end end private + def build_success_result(event, selected_json_schema_version, type_with_skipped_validation:) + BuildResult.success( + build_all_operations_for(event, record_preparer_factory.for_json_schema_version(selected_json_schema_version)), + type_with_skipped_validation: type_with_skipped_validation + ) + end + + # Builds the operations for an event whose per-record validation we skipped. + # + # Skipping validation means malformed data the schema walk would have rejected surfaces instead as + # an exception while we build the event's operations, and there is no bounded list of error types to + # enumerate. So we rescue anything and then run the validation we skipped: the validator tells us + # whether the data was actually bad, and if it was, hands the caller the same pinpointed message it + # would have gotten had we validated up front. A clean bill of health from the validator means the + # error was never about the data (a schema artifact defect, or a bug) and must not be swallowed. + def build_success_result_isolating_malformed_records(event, graphql_type_name, selected_json_schema_version) + build_success_result(event, selected_json_schema_version, type_with_skipped_validation: graphql_type_name) + rescue => exception + failed_result = validate_record_returning_failure(event, graphql_type_name, selected_json_schema_version) + # `raise` is overridden below to stop this class from *originating* an error instead of returning a + # `BuildResult`. Here we propagate one that already escaped a collaborator, which is exactly what + # happens without this rescue, so we deliberately bypass that guard. + ::Kernel.raise(exception) unless failed_result + failed_result + end + def select_json_schema_version(event) available_json_schema_versions = schema_artifacts.available_json_schema_versions @@ -117,9 +148,8 @@ def prepare_event(event) event.merge("record" => event["record"].merge("id" => event.fetch("id"))) end - def validate_record_returning_failure(event, selected_json_schema_version) + def validate_record_returning_failure(event, graphql_type_name, selected_json_schema_version) record = event.fetch("record") - graphql_type_name = event.fetch("type") validator = validator(graphql_type_name, selected_json_schema_version) if (error_message = validator.validate_with_error_message(record)) @@ -127,13 +157,48 @@ def validate_record_returning_failure(event, selected_json_schema_version) end end + # `Zlib.crc32` returns a value in `[0, 2**32)`. Pre-dividing that space by 100 lets us test a + # configured percent with a single multiply instead of dividing on every event. + CRC32_SPACE_PER_PERCENT = (1 << 32) / 100.0 + + # Decides whether to skip per-record validation for `event` of `type`. The decision is + # deterministic per event id: a stable `Zlib.crc32` of `EventID#to_s` maps each event to a + # point in the CRC32 space, and we skip validation for the configured percentage of that + # space. Same event id => same decision across pods and retries, so retries never flip a + # record between validated and skipped. `String#hash` is unsuitable here, as `RUBY_HASH_SEED` + # is per-process. The `<= 0` and `>= 100` guards keep the endpoints exact, so no float + # boundary error can make a `0` percent skip a record or a `100` percent validate one. + def skip_validation?(type, event) + percent = skip_record_validation_percents_by_type[type] + return false if percent.nil? || percent <= 0 + return true if percent >= 100 + ::Zlib.crc32(EventID.from_event(event).to_s) < percent * CRC32_SPACE_PER_PERCENT + end + def build_failed_result(event, payload_description, validation_message) message = "Malformed #{payload_description}. #{validation_message}" # Here we use the `RecordPreparer::Identity` record preparer because we may not have a valid JSON schema # version number in this case (which is usually required to get a `RecordPreparer` from the factory), and # we won't wind up using the record preparer for real on these operations, anyway. - operations = build_all_operations_for(event, RecordPreparer::Identity) + # + # Building operations for an event we already know is malformed can itself fail--for example, when the + # record omits a field an update target derives its id from. Reporting what was malformed matters more + # than reporting the operations we would have run, and `FailedEventError#operations` is documented to + # sometimes be empty for exactly this reason, so we fall back to no operations rather than let a second + # failure mask the first. + operations = begin + build_all_operations_for(event, RecordPreparer::Identity) + rescue => exception + logger.warn({ + "message_type" => "FailedEventOperationBuildingFailure", + "message_id" => event["message_id"], + "event_id" => EventID.from_event(event).to_s, + "error_class" => exception.class.name, + "error_message" => exception.message + }) + [] # : ::Array[_Operation] + end BuildResult.failure(FailedEventError.new(event: event, operations: operations.to_set, main_message: message)) end @@ -192,14 +257,17 @@ def raise(*args) # Return value from `build` that indicates what happened. # - If it was successful, `operations` will be a non-empty array of operations and `failed_event_error` will be nil. # - If there was a validation issue, `operations` will be an empty array and `failed_event_error` will be non-nil. - BuildResult = ::Data.define(:operations, :failed_event_error) do + # - `type_with_skipped_validation` names the event's GraphQL type when per-record validation was skipped + # (via `skip_record_validation_percents_by_type`), and is nil otherwise. `Processor` aggregates this + # for observability. + BuildResult = ::Data.define(:operations, :failed_event_error, :type_with_skipped_validation) do # @implements BuildResult - def self.success(operations) - new(operations, nil) + def self.success(operations, type_with_skipped_validation: nil) + new(operations, nil, type_with_skipped_validation) end def self.failure(failed_event_error) - new([], failed_event_error) + new([], failed_event_error, nil) end end end diff --git a/elasticgraph-indexer/lib/elastic_graph/indexer/processor.rb b/elasticgraph-indexer/lib/elastic_graph/indexer/processor.rb index 4f8064c39..23cfe0c3e 100644 --- a/elasticgraph-indexer/lib/elastic_graph/indexer/processor.rb +++ b/elasticgraph-indexer/lib/elastic_graph/indexer/processor.rb @@ -45,6 +45,8 @@ def process_returning_failures(events, refresh_indices: false) factory_results = factory_results_by_event.values + log_skipped_record_validations(factory_results) + bulk_result = @datastore_router.bulk(factory_results.flat_map(&:operations), refresh: refresh_indices) successful_operations = bulk_result.successful_operations(check_failures: false) @@ -62,6 +64,26 @@ def process_returning_failures(events, refresh_indices: false) private + # Emits a single aggregate log line per batch when any records had their per-record validation + # skipped (via `skip_record_validation_percents_by_type`). Skipping a safety check should never be + # silent, but per-record logging would be untenable at backfill scale (a `100` percent is one line + # per record), so we tally by type and log once. Mirrors the batch-level + # `ElasticGraphIndexingLatencies` log. + # + # A skipped record that then failed is deliberately absent from the tally: `Operation::Factory` + # re-runs the skipped validation when building its operations raises, so such a record ended up + # validated after all, and is already reported as a failed event. + def log_skipped_record_validations(factory_results) + counts_by_type = factory_results.filter_map(&:type_with_skipped_validation).tally + return if counts_by_type.empty? + + @logger.info({ + "message_type" => "RecordValidationSkipped", + "count" => counts_by_type.values.sum, + "counts_by_type" => counts_by_type + }) + end + def categorize_failures(failures, events) source_event_versions_by_cluster_by_op = @datastore_router.source_event_versions_in_index( failures.flat_map { |f| f.versioned_operations.to_a } diff --git a/elasticgraph-indexer/sig/elastic_graph/indexer/config.rbs b/elasticgraph-indexer/sig/elastic_graph/indexer/config.rbs index a4f343d59..c604057db 100644 --- a/elasticgraph-indexer/sig/elastic_graph/indexer/config.rbs +++ b/elasticgraph-indexer/sig/elastic_graph/indexer/config.rbs @@ -5,16 +5,19 @@ module ElasticGraph attr_reader latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer] attr_reader skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]] + attr_reader skip_record_validation_percents_by_type: ::Hash[::String, ::Float] attr_reader extension_modules: ::Array[::Module] def initialize: ( ?latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer], ?skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]], + ?skip_record_validation_percents_by_type: ::Hash[::String, ::Float], ?extension_modules: ::Array[::Module]) -> void def with: ( ?latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer], ?skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]], + ?skip_record_validation_percents_by_type: ::Hash[::String, ::Float], ?extension_modules: ::Array[::Module]) -> Config def self.members: () -> ::Array[::Symbol] @@ -26,6 +29,7 @@ module ElasticGraph def convert_values: ( latency_slo_thresholds_by_timestamp_in_ms: untyped, skip_derived_indexing_type_updates: untyped, + skip_record_validation_percents_by_type: untyped, extension_modules: untyped ) -> ::Hash[::Symbol, untyped] diff --git a/elasticgraph-indexer/sig/elastic_graph/indexer/operation/factory.rbs b/elasticgraph-indexer/sig/elastic_graph/indexer/operation/factory.rbs index 82e968053..1b27e94c7 100644 --- a/elasticgraph-indexer/sig/elastic_graph/indexer/operation/factory.rbs +++ b/elasticgraph-indexer/sig/elastic_graph/indexer/operation/factory.rbs @@ -11,6 +11,7 @@ module ElasticGraph attr_reader record_preparer_factory: RecordPreparer::Factory attr_reader logger: ::Logger attr_reader skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]] + attr_reader skip_record_validation_percents_by_type: ::Hash[::String, ::Float] attr_reader configure_record_validator: (^(validatorFactory) -> validatorFactory)? def initialize: ( @@ -19,6 +20,7 @@ module ElasticGraph record_preparer_factory: RecordPreparer::Factory, logger: ::Logger, skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]], + skip_record_validation_percents_by_type: ::Hash[::String, ::Float], configure_record_validator: (^(validatorFactory) -> validatorFactory)? ) -> void @@ -28,6 +30,7 @@ module ElasticGraph ?record_preparer_factory: RecordPreparer::Factory, ?logger: ::Logger, ?skip_derived_indexing_type_updates: ::Hash[::String, ::Set[::String]], + ?skip_record_validation_percents_by_type: ::Hash[::String, ::Float], ?configure_record_validator: (^(validatorFactory) -> validatorFactory)? ) -> instance end @@ -44,7 +47,11 @@ module ElasticGraph def select_json_schema_version: (event) { (BuildResult) -> bot } -> (::Integer | bot) def prepare_event: (event) -> event - def validate_record_returning_failure: (event, ::Integer) -> BuildResult? + def validate_record_returning_failure: (event, ::String, ::Integer) -> BuildResult? + def build_success_result: (event, ::Integer, type_with_skipped_validation: ::String?) -> BuildResult + def build_success_result_isolating_malformed_records: (event, ::String, ::Integer) -> BuildResult + CRC32_SPACE_PER_PERCENT: ::Float + def skip_validation?: (::String, event) -> bool def build_failed_result: (event, ::String, ::String) -> BuildResult def build_all_operations_for: (event, _RecordPreparer) -> ::Array[_Operation] def index_definitions_for: (::String) -> ::Array[DatastoreCore::_IndexDefinition] @@ -53,15 +60,17 @@ module ElasticGraph class BuildResult attr_reader operations: ::Array[_Operation] attr_reader failed_event_error: FailedEventError? + attr_reader type_with_skipped_validation: ::String? - def initialize: (::Array[_Operation], FailedEventError?) -> void + def initialize: (::Array[_Operation], FailedEventError?, ::String?) -> void def with: ( ?operations: ::Array[_Operation], - ?failed_event_error: FailedEventError? + ?failed_event_error: FailedEventError?, + ?type_with_skipped_validation: ::String? ) -> BuildResult - def self.success: (::Array[_Operation]) -> BuildResult + def self.success: (::Array[_Operation], ?type_with_skipped_validation: ::String?) -> BuildResult def self.failure: (FailedEventError) -> BuildResult end end diff --git a/elasticgraph-indexer/sig/elastic_graph/indexer/processor.rbs b/elasticgraph-indexer/sig/elastic_graph/indexer/processor.rbs index ac450f43a..0f34056ce 100644 --- a/elasticgraph-indexer/sig/elastic_graph/indexer/processor.rbs +++ b/elasticgraph-indexer/sig/elastic_graph/indexer/processor.rbs @@ -20,6 +20,7 @@ module ElasticGraph @indexing_latency_slo_thresholds_by_timestamp_in_ms: ::Hash[::String, ::Integer] @clock: singleton(::Time) + def log_skipped_record_validations: (::Array[Operation::Factory::BuildResult]) -> void def categorize_failures: (::Array[FailedEventError], ::Array[event]) -> ::Array[FailedEventError] def calculate_latency_metrics: (::Array[_Operation], ::Array[Operation::Result]) -> void end diff --git a/elasticgraph-indexer/spec/unit/elastic_graph/indexer/config_spec.rb b/elasticgraph-indexer/spec/unit/elastic_graph/indexer/config_spec.rb index 2c157a8e9..41fea0a71 100644 --- a/elasticgraph-indexer/spec/unit/elastic_graph/indexer/config_spec.rb +++ b/elasticgraph-indexer/spec/unit/elastic_graph/indexer/config_spec.rb @@ -32,6 +32,33 @@ class Indexer expect(config.skip_derived_indexing_type_updates).to eq("WidgetCurrency" => ["USD"].to_set) end + it "coerces `skip_record_validation_percents_by_type` percents to floats (so integer YAML values like `90` become `90.0`)" do + config = Config.from_parsed_yaml("indexer" => { + "latency_slo_thresholds_by_timestamp_in_ms" => {}, + "skip_record_validation_percents_by_type" => { + "Widget" => 90, + "Component" => 99.5 + } + }) + + expect(config.skip_record_validation_percents_by_type).to eq("Widget" => 90.0, "Component" => 99.5) + expect(config.skip_record_validation_percents_by_type.values).to all(be_a(::Float)) + end + + it "rejects `skip_record_validation_percents_by_type` percents outside `[0, 100]`" do + expect { + Config.from_parsed_yaml("indexer" => { + "skip_record_validation_percents_by_type" => {"Widget" => 100.5} + }) + }.to raise_error Errors::ConfigError + + expect { + Config.from_parsed_yaml("indexer" => { + "skip_record_validation_percents_by_type" => {"Widget" => -0.1} + }) + }.to raise_error Errors::ConfigError + end + describe "#extension_modules", :in_temp_dir do it "loads the extension modules from disk" do File.write("eg_extension_module1.rb", <<~EOS) diff --git a/elasticgraph-indexer/spec/unit/elastic_graph/indexer/operation/factory_spec.rb b/elasticgraph-indexer/spec/unit/elastic_graph/indexer/operation/factory_spec.rb index 2b3b5df10..e2162f611 100644 --- a/elasticgraph-indexer/spec/unit/elastic_graph/indexer/operation/factory_spec.rb +++ b/elasticgraph-indexer/spec/unit/elastic_graph/indexer/operation/factory_spec.rb @@ -75,6 +75,250 @@ module Operation end end + # We deliberately construct the indexer here without going through `build_indexer`. The + # `skip_record_validation_percents_by_type` knob is intentionally not exposed via spec helpers so that + # tests cannot silently weaken validation; enabling it must be a visible, deliberate choice + # in each spec that exercises it. + context "when the indexer is configured to skip record validation for some types" do + let(:indexer) do + datastore_core = build_datastore_core + Indexer.new( + datastore_core: datastore_core, + config: Indexer::Config.new( + latency_slo_thresholds_by_timestamp_in_ms: {}, + skip_derived_indexing_type_updates: {}, + skip_record_validation_percents_by_type: {"Component" => 100} + ) + ) + end + + it "skips per-type record validation for the listed type but still builds operations" do + event = build_upsert_event(:component, id: "1", __version: 1) + event["record"]["name"] = 123 # would normally fail JSON schema validation + + expect(build_expecting_success(event)).to eq([new_primary_indexing_operation({ + "op" => "upsert", + "id" => "1", + "type" => "Component", + "version" => 1, + "record" => event["record"], + JSON_SCHEMA_VERSION_KEY => 1 + })]) + end + + it "records the skipped type on the successful build result" do + event = build_upsert_event(:component, id: "1", __version: 1) + + result = indexer.operation_factory.build(event) + + expect(result.type_with_skipped_validation).to eq("Component") + end + + it "still validates record-level fields for types that are not in the skip list" do + widget_event = build_upsert_event(:widget, id: "1", __version: 1) + widget_event["record"]["name"] = 123 + + expect_failed_event_error(widget_event, "Malformed Widget record", "name") + end + + it "does not flag validated types on the successful build result" do + event = build_upsert_event(:widget, id: "1", __version: 1) + + result = indexer.operation_factory.build(event) + + expect(result.type_with_skipped_validation).to be_nil + end + + it "still applies envelope-level validation for skipped types" do + event = build_upsert_event(:component, id: "1", __version: -1) + + expect_failed_event_error(event, "/properties/version") + end + end + + context "when the indexer configures a partial skip percent for a type" do + let(:indexer) do + datastore_core = build_datastore_core + Indexer.new( + datastore_core: datastore_core, + config: Indexer::Config.new( + latency_slo_thresholds_by_timestamp_in_ms: {}, + skip_derived_indexing_type_updates: {}, + skip_record_validation_percents_by_type: {"Component" => 50} + ) + ) + end + + it "skips validation when the event's point in the crc32 space falls in the skipped portion" do + # Stub crc32 to the bottom of the space -- well inside the skipped 50% -- forcing the "skip" branch. + allow(::Zlib).to receive(:crc32).and_return(0) + + event = build_upsert_event(:component, id: "1", __version: 1) + event["record"]["name"] = 123 # would normally fail JSON schema validation + + expect { + build_expecting_success(event) + }.not_to raise_error + end + + it "still validates when the event's point in the crc32 space falls outside the skipped portion" do + # Stub crc32 to 75% of the way through the space, past the skipped 50%, forcing validation. + allow(::Zlib).to receive(:crc32).and_return((2**32 * 0.75).to_i) + + event = build_upsert_event(:component, id: "1", __version: 1) + event["record"]["name"] = 123 + + expect_failed_event_error(event, "Malformed Component record", "name") + end + + it "produces the same skip decision for the same event on retry" do + # No stubbing -- this exercises the real crc32 and locks in determinism without + # coupling to a specific hash output. Two builds of the same event must agree on + # whether to skip validation. + event = build_upsert_event(:component, id: "1", __version: 1) + event["record"]["name"] = 123 # would fail validation if not skipped + + first = indexer.operation_factory.build(event) + second = indexer.operation_factory.build(event) + + expect(first.failed_event_error.nil?).to eq(second.failed_event_error.nil?) + expect(first.operations.size).to eq(second.operations.size) + expect(first.type_with_skipped_validation).to eq(second.type_with_skipped_validation) + end + end + + context "when record validation is skipped for a type that has derived-index update targets" do + # Widget has a `WidgetCurrency` derived index update target. With validation skipped, + # `build_all_operations_for` still has to traverse `Update.operations_for` and the + # schema artifacts. This spec locks in that skipping does not regress the derived path. + let(:indexer) do + datastore_core = build_datastore_core + Indexer.new( + datastore_core: datastore_core, + config: Indexer::Config.new( + latency_slo_thresholds_by_timestamp_in_ms: {}, + skip_derived_indexing_type_updates: {}, + skip_record_validation_percents_by_type: {"Widget" => 100} + ) + ) + end + + it "still emits both the primary and the derived-index update operations" do + event = build_upsert_event(:widget, id: "1", __version: 1) + formatted_event = { + "op" => "upsert", + "id" => "1", + "type" => "Widget", + "version" => 1, + "record" => event["record"], + JSON_SCHEMA_VERSION_KEY => 1 + } + + expect(build_expecting_success(event)).to contain_exactly( + new_primary_indexing_operation(formatted_event, index_def: index_def_named("widgets")), + widget_currency_derived_update_operation_for(formatted_event) + ) + end + end + + context "when building the operations for a record whose validation was skipped raises" do + # These specs cover the re-validation gate: an exception escaping operation building is + # answered by running the validation we skipped. If the validator faults the record, the + # caller gets the same failure it would have gotten had we validated up front. If the + # validator is happy, the error was never about the data and must not be swallowed. + let(:indexer) do + datastore_core = build_datastore_core + Indexer.new( + datastore_core: datastore_core, + config: Indexer::Config.new( + latency_slo_thresholds_by_timestamp_in_ms: {}, + skip_derived_indexing_type_updates: {}, + skip_record_validation_percents_by_type: {"Widget" => 100} + ) + ) + end + + it "reports a value the indexing preparers can't coerce as a failed event carrying the validator's message" do + event = build_upsert_event(:widget, id: "1", __version: 1) + # `IndexingPreparers::Integer` raises `Errors::IndexOperationError` on this; validation + # would have caught it first as a type mismatch on `amount_cents`. + event["record"]["cost"] = {"currency" => "USD", "amount_cents" => "not a number"} + + message = expect_failed_event_error(event, "Malformed Widget record", "amount_cents") + + expect(message).to exclude("IndexOperationError") + end + + it "reports a missing or unknown abstract-type `__typename` as a failed event carrying the validator's message" do + # Widget's `inventor` field is the `Inventor` union, whose JSON schema requires + # `__typename` and pins a `const` discriminator on each concrete subtype. Both the + # missing and the unknown case therefore fail validation, which is why no dedicated + # error type is needed for them. + event = build_upsert_event(:widget, id: "1", __version: 1) + event["record"]["inventor"] = {"__typename" => "NotARealConcreteType", "name" => "anon"} + + expect_failed_event_error(event, "Malformed Widget record", "inventor") + end + + it "re-raises an error the validator has no opinion about, so bugs are not hidden as data failures" do + event = build_upsert_event(:widget, id: "1", __version: 1) + + # Asserting on class *and* message matters here: `raise` is overridden on + # `Operation::Factory` to reject raising from the class, so a bare `raise exception` + # would substitute the guard's own error and still satisfy a bare `raise_error`. + expect { + factory_whose_record_preparation_is_broken.build(event) + }.to raise_error(::KeyError, a_string_including("nameInIndex")) + end + end + + it "does not rescue when validation was not skipped, so the validated path behaves exactly as before" do + event = build_upsert_event(:widget, id: "1", __version: 1) + + expect { + factory_whose_record_preparation_is_broken.build(event) + }.to raise_error(::KeyError, a_string_including("nameInIndex")) + end + + context "when an event is malformed in a way that also breaks building its operations", :expect_warning_logging do + # `Widget` requires `cost`, and its derived `WidgetCurrency` update target sources its id + # from `cost.currency`, so a `Widget` with no `cost` both fails validation and breaks + # building the operations we attach to the `FailedEventError`. Reporting the malformation + # matters more than reporting operations we are never going to run, and + # `FailedEventError#operations` is documented as sometimes being empty for this reason. + let(:event) do + build_upsert_event(:widget, id: "1", __version: 1).tap { |e| e["record"].delete("cost") } + end + + it "still reports what was malformed instead of letting the second failure mask the first" do + failure = indexer.operation_factory.build(event).failed_event_error + + expect(failure).to be_an(FailedEventError) + expect(failure.operations).to be_empty + expect(failure.main_message).to include("Malformed Widget record", "cost").and exclude("Key not found") + end + + it "logs the failure it swallowed, so a discarded operation-building error is still traceable" do + indexer.operation_factory.build(event) + + expect(logged_jsons_of_type("FailedEventOperationBuildingFailure")).to match([a_hash_including( + "event_id" => "Widget:1@v1", + "error_class" => "KeyError" + )]) + end + end + + # A factory whose record preparation fails for a reason record validation cannot detect: a + # runtime metadata defect, standing in for any bug that is not about the data itself. + def factory_whose_record_preparation_is_broken + record_preparer_factory = instance_double(RecordPreparer::Factory) + + allow(record_preparer_factory).to receive(:for_json_schema_version) + .and_raise(::KeyError, 'key not found: "nameInIndex"') + + indexer.operation_factory.with(record_preparer_factory: record_preparer_factory) + end + it "generates a primary indexing operation for a single index with latency metrics" do event = build_upsert_event(:component, id: "1", __version: 1) latency_timestamps = {"latency_timestamps" => {"created_in_esperanto_at" => "2012-04-23T18:25:43.511Z"}} diff --git a/elasticgraph-indexer/spec/unit/elastic_graph/indexer/processor_spec.rb b/elasticgraph-indexer/spec/unit/elastic_graph/indexer/processor_spec.rb index 79a47d7da..ad5e941f4 100644 --- a/elasticgraph-indexer/spec/unit/elastic_graph/indexer/processor_spec.rb +++ b/elasticgraph-indexer/spec/unit/elastic_graph/indexer/processor_spec.rb @@ -386,6 +386,47 @@ def make_component_bad(component) end end + context "when the indexer skips per-record validation for some types" do + # `skip_record_validation_percents_by_type` is intentionally not exposed via `build_indexer`, so we + # construct the indexer directly (reusing the spec's router spy and clock) to enable it. + let(:indexer) do + Indexer.new( + datastore_core: build_datastore_core, + config: Indexer::Config.new( + latency_slo_thresholds_by_timestamp_in_ms: {}, + skip_derived_indexing_type_updates: {}, + skip_record_validation_percents_by_type: {"Component" => 100} + ), + datastore_router: datastore_router, + clock: clock + ) + end + + it "logs a single aggregate `RecordValidationSkipped` entry per batch, counted by type" do + component1 = build_upsert_event(:component, id: "c1", __version: 1) + component2 = build_upsert_event(:component, id: "c2", __version: 1) + address = build_upsert_event(:address, id: "a1", __version: 1) + + process([component1, component2, address]) + + expect(logged_jsons_of_type("RecordValidationSkipped")).to contain_exactly( + a_hash_including( + "message_type" => "RecordValidationSkipped", + "count" => 2, + "counts_by_type" => {"Component" => 2} + ) + ) + end + + it "logs nothing when no records in the batch had validation skipped" do + address = build_upsert_event(:address, id: "a1", __version: 1) + + process([address]) + + expect(logged_jsons_of_type("RecordValidationSkipped")).to be_empty + end + end + def build_indexer_with(latency_thresholds:) build_indexer( clock: clock, diff --git a/elasticgraph-local/lib/elastic_graph/local/spec_support/config_schema.yaml b/elasticgraph-local/lib/elastic_graph/local/spec_support/config_schema.yaml index 638a74a12..9b38c9251 100644 --- a/elasticgraph-local/lib/elastic_graph/local/spec_support/config_schema.yaml +++ b/elasticgraph-local/lib/elastic_graph/local/spec_support/config_schema.yaml @@ -508,6 +508,25 @@ properties: - {} - WidgetWorkspace: - ABC12345678 + skip_record_validation_percents_by_type: + description: |- + Map of GraphQL type names to the percentage of records of that type whose per-record JSON schema validation should be skipped. `0` (or an absent key) validates every record of the type; `100` skips every record; values in between sample, and may be fractional. The decision is deterministic per event id (`type:id@vversion`), so the same event makes the same choice on every retry and on every indexer pod. The event envelope (op, id, type, version, json_schema_version, latency_timestamps) is always validated, regardless of this setting. + + With a large schema the per-record schema walk consumes a significant share of indexing CPU: every record is checked against every regex, enum, min/max, format, and abstract-type discriminator defined for its type. Skipping it trades that check for throughput, which is worthwhile when backfilling data that was already validated upstream. Leaving a percentage of records validated keeps a canary in place so schema drift still surfaces. + + Note: skipping validation makes malformed-data detection later and less precise. A malformation found while building an event's operations is still reported as an isolated event failure, carrying the message validation itself would have produced. But one found only while serializing an operation for the datastore (an unparsable rollover index timestamp, or a missing custom routing field) raises an error that fails the entire batch, including the well-formed events in it. Since such a batch produces no partial-failure response, the queue redelivers all of its events, and the malformed record fails them again on each retry until it is drained to the dead letter queue. Leave this empty for live-traffic ingestion. + type: object + patternProperties: + "^[A-Z]\\w*$": + type: number + minimum: 0 + maximum: 100 + additionalProperties: false + default: {} + examples: + - {} + - Widget: 90 + Component: 100 extension_modules: description: Array of modules that will be extended onto the component instance to support extension libraries.