Repository navigation
fix(memory): replicate LTM updates and deletes and make creates retry-safe in the multi-Region replication sample - #2172
Open
gilbertlep wants to merge 3 commits into
Conversation
added 3 commits
October 9, 2026 12:07
…-safe in the multi-Region replication sample - Creates send a clientToken derived from the stream event, so a redelivered event can't create a duplicate (requestIdentifier is for tracking only). - A small record map (DynamoDB in the stack, in-memory for the demo) remembers which target record each source record became, so updates go to BatchUpdateMemoryRecords and deletes to BatchDeleteMemoryRecords, and stale or redelivered events are skipped. - README, docstrings, and the demo message no longer claim that requestIdentifier makes replays idempotent or that consolidation handles deletes. - failedRecords[].errorCode is an integer in the API, but the sample compared it to exception names, so a throttled record was dropped as terminal. 429 and 5xx now raise so the batch is retried; 404 on a delete counts as done. - Event times from the stream have 9 fractional digits, which datetime.fromisoformat rejects before Python 3.11; they are trimmed to 6, so ordering works on the README's minimum Python (3.10). - dual_writer: document that retries must reuse client_token and event_timestamp. - Unit tests: python -m pytest tests (no AWS access needed). - Changed files pass ruff 0.16.10 (latest, as CI installs it); this includes small mechanical fixes to code that was already there. - Add gilbertlep to CONTRIBUTORS.md.
…streaming in deploy_streaming.sh Setting the execution role leaves the memory UPDATING for a minute or two, so the next UpdateMemory (enable_streaming.py) failed with "Memory is in transitional state UPDATING". full_demo.py already waits at the same point.
ListMemoryRecords with namespace="/" returns no records, so the demo reported 0 source LTM records, waited out both timeouts, and couldn't show LTM replication. namespacePath="/" lists every namespace (available since the sample's minimum boto3, 1.43.36).
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Issue number: Fixes #2171
Concise description of the PR
Changes the LTM stream consumer in
00-multi-region-replicationso the target memory stays in sync with the source, because updates became duplicates, deletes were skipped, Lambda retries could duplicate creates, and throttled records were dropped (details in the issue).BatchCreateMemoryRecordsBatchUpdateMemoryRecordson that record.BatchDeleteMemoryRecords, and the map keeps a tombstone (a "deleted" marker), so a late create or update for that record is skipped.clientTokenclientTokenis a hash of the stream event (memory ID, record ID, event type, event time), so a redelivered event sends the same token.errorCodecompared to stringsClientErrorso the batch is retried. 404 on a delete counts as done. Anything else stays terminal for that record.eventTimehas 9 fractional digitsdatetime.fromisoformatrejects 9 before Python 3.11, and the README's minimum is 3.10, so ordering silently used "now" there.deploy_streaming.shfails on its last stepACTIVEafter setting the execution role, asfull_demo.pyalready does. Separate commit.full_demo.pyreports 0 LTM recordsnamespacePath="/"(in boto3 since the sample's minimum, 1.43.36). Separate commit.The record map is a
RecordMapprotocol with two implementations:DynamoDBRecordMap(Lambda path; newRecordMapTableininfra/streaming-stack.yaml, on-demand, PITR and SSE on) andInMemoryRecordMap(local demo and tests). No new Python dependencies.Also:
requestIdentifiermakes replays idempotent, that consolidation handles deletes, or that a reverse pass lands only what's missing. The Deletes limitation is replaced with the one real gap: records copied without a map entry (for example, an unrecorded backfill).dual_writer.py: docstring now says that a retry must reuseclient_tokenandevent_timestamp.Optional[X]toX | None,dict()to literals, import order).User experience
Before: after an update or delete in the source Region, the target returns stale or deleted records, and retries can leave duplicates.
After: the target applies creates, updates, and deletes once each, in order, including when Lambda retries a batch.
Verification
Live, us-east-1 to us-west-2, through the deployed stack (
deploy_streaming.sh, then source change, Kinesis, Lambda, target). Run twice, the second time on the final code:BatchCreateMemoryRecords)skipped: 1; still one target recordclientTokenresent)replicated: 1; the API returned the same target record ID; still one recordskipped: 1; target stays emptyLambda logs over both runs: 16 invocations, 8 replicated, 8 skipped, 0 failed, 0 errors.
Direct API probes on the target memory, which the error handling relies on:
BatchCreateMemoryRecordswith the sameclientTokenreturns the originalsuccessfulRecords(samememoryRecordId) and creates nothing new.failedRecordswith"errorCode": 404, "errorMessage": "Memory record not found". 429 and 5xx weren't observed; they follow the same HTTP-status convention.deploy_streaming.shfailed on its last step before the second commit and completed after it.scripts/full_demo.py(the README's headline demo): before the third commit it reported 0 LTM records on both sides while replicating 5. After it: 5 source records extracted, 5 replicated (received=6 replicated=5 skipped=1 failed=0), the target holds exactly 5 with no duplicates, STM 8 of 8 events, and teardown completed.Offline:
python -m pytest testsfrom the sample directory, 15 passed on Python 3.10, 3.12, and 3.13, no AWS access needed. They cover redelivery, update, stale update, unmapped update, delete, create after delete, delete of a missing record, control events, retryable vs terminal failures, nanosecond event times, and the DynamoDB tombstone round trip.ruff checkandruff format --checkclean at 0.16.10.Checklist
Acknowledgment
By submitting this pull request, I confirm that you can use, modify, copy, and redistribute this contribution, under the terms of the project license.