Skip to content

Preserve supersession for malformed envelopes - #1404

Closed
jwils wants to merge 1 commit into
mainfrom
joshuaw/preserve-envelope-supersession
Closed

jwils wants to merge 1 commit into
mainfrom
joshuaw/preserve-envelope-supersession

Conversation

@jwils

@jwils jwils commented Sep 21, 2026

Copy link
Copy Markdown
Collaborator

Why

After #1397 validates envelopes before indexing, malformed envelopes no longer reach the supersession check. They keep failing on retry even after a newer event replaces them in the datastore.

What

Preserve supersession for rejected JSON envelopes whose operation, type, ID, and version can still be validated. Keep the original failure when identity or destination targets cannot be established.

How

Build read-only version lookups and run them through the existing supersession check after valid events are indexed. Every target must have a strictly newer stored version. Use the validated envelope ID when the record contains a conflicting ID. Rejected payloads never become indexable Event values.

Risk

This changes which malformed messages are retried. Equal or missing stored versions remain failures, and rejected payloads never produce datastore writes.

Testing

No manual testing. Regression specs cover newer versus equal versions, conflicting record IDs, missing records, invalid identities, and unavailable targets.

Bigger picture

Follow-up to #1398. This PR contains only supersession handling; JSON API changes and non-object JSON handling are separate work.

@jwils
jwils added this pull request to stack #1401 September 21, 2026 02:14
@myronmarston
myronmarston force-pushed the joshuaw/preserve-envelope-supersession branch from 8501069 to f2689b4 Compare September 21, 2026 02:30
@jwils
jwils marked this pull request as ready for review September 22, 2026 16:16
Base automatically changed from myron/typed-event-object to main September 22, 2026 18:27
@myronmarston
myronmarston force-pushed the joshuaw/preserve-envelope-supersession branch from f2689b4 to 6ff7ad0 Compare September 22, 2026 18:27

@myronmarston myronmarston left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I started reviewing this and found some things to request changes on but I think I'm in favor of not fixing this issue and leaving it as is.

If you ask me in isolation: "should ElasticGraph malformed even supersession work for events with malformed envelopes?" I'd say "yes"...but this adds a lot of complexity for a case I don't expect anyone to ever hit.

  • JSONIngestion::Indexer#process_returning_failures gets way more complicated...and presumably obligates ProtoIngestion::Indexer#process_returning_failures to also take on that kind of complexity for parity. I'm particularly sensitive to adding that kind of complexity to each ingestion adapter.
  • An event that's malformed in the event envelope rather than the record is really malformed. I think it's OK if we say we aren't going to handle those ones as smoothly.
  • In practice, I don't think we've ever run into malformed data in the event envelope. It's the kind of thing that would likely only happen when you first publish into EG for the first time before you're fully in prod.
  • The "has this event been superseded?" check inside ElasticGraph can be pretty inefficient (see #1077 and #1078 for details--it caused us a load issue at one point!). An issue with the envelope is likely to be widespread rather than isolated to a couple of events and preserving supersession for them could take down the cluster.

Putting these all together, my instinct is to close this as a "won't fix".

Thoughts?

# @!attribute [r] destination_index_def
# @return [DatastoreCore::IndexDefinition] the index to search
# @!attribute [r] update_target
# @return [SchemaArtifacts::RuntimeMetadata::UpdateTarget] the relationship whose version to read

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
# @return [SchemaArtifacts::RuntimeMetadata::UpdateTarget] the relationship whose version to read
# @return [SchemaArtifacts::RuntimeMetadata::UpdateTarget] the target of the update (from which the version will be read)


# @return [Boolean] whether this target stores source event versions
def versioned?
update_target.for_normal_indexing?

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
update_target.for_normal_indexing?
# We do not track source event versions when applying derived indexing updates, but we do for
# normal indexing updates, so if the update target is for normal indexing it's a versioned operation.
update_target.for_normal_indexing?

I couldn't remember why versioned? == update_target.for_normal_indexing? and found this that explained it:

def versioned?
# We do not track source event versions when applying derived indexing updates, but we do for
# normal indexing updates, so if the update target is for normal indexing it's a versioned operation.
update_target.for_normal_indexing?
end

Might as well carry it forward.

end

ElasticGraph::Indexer::EventID.new(type: identity.fetch("type"), id: identity.fetch("id"), version: identity.fetch("version").to_i)
end

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

diff --git a/elasticgraph-indexer/lib/elastic_graph/indexer/event_id.rb b/elasticgraph-indexer/lib/elastic_graph/indexer/event_id.rb
index 3b98991f..3492fa9f 100644
--- a/elasticgraph-indexer/lib/elastic_graph/indexer/event_id.rb
+++ b/elasticgraph-indexer/lib/elastic_graph/indexer/event_id.rb
@@ -22,7 +22,7 @@ module ElasticGraph
       # @param hash [Hash<String, Object>] a decoded indexing payload
       # @return [EventID]
       def self.from_decoded_hash(hash)
-        new(type: hash["type"], id: hash["id"], version: hash["version"])
+        new(type: hash["type"], id: hash["id"], version: hash["version"]&.to_i)
       end
 
       def to_s
diff --git a/elasticgraph-json_ingestion/lib/elastic_graph/json_ingestion/envelope_validator.rb b/elasticgraph-json_ingestion/lib/elastic_graph/json_ingestion/envelope_validator.rb
index ffe2eea4..3c04100c 100644
--- a/elasticgraph-json_ingestion/lib/elastic_graph/json_ingestion/envelope_validator.rb
+++ b/elasticgraph-json_ingestion/lib/elastic_graph/json_ingestion/envelope_validator.rb
@@ -53,23 +53,25 @@ module ElasticGraph
       # @param decoded_event [Hash<String, Object>] a decoded JSON event that failed envelope validation
       # @return [ElasticGraph::Indexer::EventID, nil] validated identity, if available
       def identity_for_supersession(decoded_event)
-        requested_version = decoded_event[JSON_SCHEMA_VERSION_KEY]
-        version = if requested_version.is_a?(::Integer)
-          closest_available_json_schema_version(requested_version)
+        requested_json_schema_version = decoded_event[JSON_SCHEMA_VERSION_KEY]
+        json_schema_version = if requested_json_schema_version.is_a?(::Integer)
+          closest_available_json_schema_version(requested_json_schema_version)
         else
           @schema_artifacts.available_json_schema_versions.max
         end
-        return unless version
+        return unless json_schema_version
 
-        envelope_schema = validator(EVENT_ENVELOPE_JSON_SCHEMA_NAME, version).schema
-        identity = decoded_event.slice("op", "type", "id", "version")
-        return unless %w[op type id version].all? do |field|
-          identity.key?(field) && envelope_schema.ref("#/$defs/#{EVENT_ENVELOPE_JSON_SCHEMA_NAME}/properties/#{field}").valid?(identity.fetch(field))
+        envelope_schema = validator(EVENT_ENVELOPE_JSON_SCHEMA_NAME, json_schema_version).schema
+        identity = NIL_EVENT.merge(decoded_event).slice(*NIL_EVENT.keys)
+        return unless identity.all? do |field, value|
+          envelope_schema.ref("#/$defs/#{EVENT_ENVELOPE_JSON_SCHEMA_NAME}/properties/#{field}").valid?(value)
         end
 
-        ElasticGraph::Indexer::EventID.new(type: identity.fetch("type"), id: identity.fetch("id"), version: identity.fetch("version").to_i)
+        ElasticGraph::Indexer::EventID.from_decoded_hash(identity)
       end
 
+      NIL_EVENT = {"op" => nil, "type" => nil, "id" => nil, "version" => nil}
+
       # The requested version might not necessarily be available (if the publisher is deployed ahead of the indexer, or an old schema
       # version is removed prematurely, or an indexer deployment is rolled back). So the behavior is to always pick the closest-available
       # version. If there's an exact match, great. Even if not an exact match, if the incoming event payload conforms to the closest match,

Some things we can improve here.

#
# @param decoded_event [Hash<String, Object>] a decoded JSON event that failed envelope validation
# @return [ElasticGraph::Indexer::EventID, nil] validated identity, if available
def identity_for_supersession(decoded_event)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Can we call this supersession_identity?

Comment on lines +50 to +54
# Extracts only the identity needed to check supersession from a malformed payload.
# This never makes a malformed event eligible to be indexed.
#
# @param decoded_event [Hash<String, Object>] a decoded JSON event that failed envelope validation
# @return [ElasticGraph::Indexer::EventID, nil] validated identity, if available

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Suggested change
# Extracts only the identity needed to check supersession from a malformed payload.
# This never makes a malformed event eligible to be indexed.
#
# @param decoded_event [Hash<String, Object>] a decoded JSON event that failed envelope validation
# @return [ElasticGraph::Indexer::EventID, nil] validated identity, if available
# Extract the identity needed to check supersession in a way that that can be used on a malformed payload.
#
# @param decoded_event [Hash<String, Object>] a decoded JSON event (e.g. that failed envelope validation)
# @return [ElasticGraph::Indexer::EventID, nil] event identity, if available

There's nothing about this method that requires the event to be invalid...it's just intended for that use but we shouldn't overly dictate how callers can use a method.

@jwils

jwils commented Sep 24, 2026

Copy link
Copy Markdown
Collaborator Author

Yep it makes a lot of sense to me. I'll close this

@jwils jwils closed this Sep 24, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants