Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,9 @@ extracted knowledge intact.
│ │ dual-write at write time │ dual-write│ ="SKIP") history only, │
│ ▼ │ │ NO re-extraction │
│ LTM extraction │ │ │
│ │ │ │ BatchCreateMemoryRecords │
│ ▼ │ │ (requestIdentifier = │
│ Kinesis record stream ─────────┼── LTM ─────┼─▶ source memoryRecordId) ◀── │
│ │ │ │ Batch Create/Update/Delete │
│ ▼ │ │ MemoryRecords, matched by │
│ Kinesis record stream ─────────┼── LTM ─────┼─▶ source record ID ◀── │
│ (MEMORY_RECORDS/FULL_CONTENT) │ streaming │ consumer Lambda │
└────────────────────────────────┘ └──────────────────────────────┘
```
Expand All @@ -36,16 +36,18 @@ Two AgentCore building blocks make this work:
records already arrive over the stream, so the target must *not* re-extract the
replicated events into duplicate records.
* **Record streaming (`streamDeliveryResources`)** — the source memory publishes
every extracted record to a Kinesis Data Stream, which the consumer replays into
the target. Using the source `memoryRecordId` as the target `requestIdentifier`
makes replays idempotent (re-delivery is a conditional no-op, not a duplicate).
every record change to a Kinesis Data Stream, which the consumer applies to the
target: creates, updates, and deletes. Creates send a `clientToken` derived from
the stream event, so re-delivery can't duplicate a record, and a DynamoDB record
map tracks which target record each source record became.

## Repository layout

```
.
├── agentcore_replication/ # reusable Python package
│ ├── stream_consumer.py # LTM: consume Kinesis record stream -> target
│ ├── record_map.py # LTM: source -> target record ID map (DynamoDB)
│ └── dual_writer.py # STM: dual-write CreateEvent (SKIP on target)
├── lambda/
│ └── stream_handler.py # Kinesis-triggered LTM consumer (production)
Expand All @@ -56,6 +58,7 @@ Two AgentCore building blocks make this work:
│ ├── enable_streaming.py # turn record streaming on/off (UpdateMemory)
│ ├── full_demo.py # end-to-end STM+LTM demo (the headline sample)
│ └── deploy_streaming.sh # deploy the LTM streaming path
├── tests/ # unit tests: python -m pytest tests
└── requirements.txt
```

Expand Down Expand Up @@ -136,21 +139,19 @@ This creates the Kinesis stream + consumer Lambda (`infra/streaming-stack.yaml`)
in the **source** region, wires the Event Source Mapping with a DLQ and CloudWatch
alarms, attaches a streaming execution role to the source memory, and enables
record streaming. From then on, extracted LTM records flow source → Kinesis →
Lambda → `BatchCreateMemoryRecords` in the target.
Lambda → the target (create, update, or delete).

> **Note:** record streaming only carries records created *after* you enable it.
> Enable streaming before your agents start writing, or backfill pre-existing data
> with `ListMemoryRecords` → `BatchCreateMemoryRecords` (the same idempotent
> `requestIdentifier` path the consumer uses).
> with `ListMemoryRecords` → `BatchCreateMemoryRecords`, and write each copy to the
> record map so that later updates and deletes find it.

## Failover

This is active-passive. On primary-region failure, point your agents at the replica
memory in the target region. To replicate the other direction afterward, swap
source/target: enable streaming on the new primary (`enable_streaming.py`), point
your dual-writer the other way, and deploy a consumer in the new source region.
Because LTM writes are idempotent, the first reverse pass safely lands only what's
missing.

> **Pause replication** during failover with
> `scripts/enable_streaming.py --memory-id <SRC> --region <SRC_REGION> --disable`.
Expand All @@ -164,9 +165,9 @@ missing.

## Limitations

* **Deletes** are not replicated (`MemoryRecordDeleted` events are skipped).
AgentCore consolidation handles stale records in the target; call
`DeleteMemoryRecord` explicitly if you need exact parity.
* **Record map coverage**: updates and deletes rely on the record map. A record
that reached the target without a map entry (for example, an unrecorded backfill)
is re-created on its next update, and its delete is skipped.
* **STM latency**: STM is replicated synchronously at write time, so a target
outage surfaces on the write path — wrap `record_turn` to tolerate target
failures (the source write is what your agent depends on).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,10 @@
* **LTM (long-term records) via record streaming** — the source memory is
configured with ``streamDeliveryResources`` (``MEMORY_RECORDS`` /
``FULL_CONTENT``), so extracted records are published to a Kinesis Data Stream.
A consumer (:mod:`stream_consumer`, run as a Lambda or locally) re-creates each
record in the target via ``BatchCreateMemoryRecords``, using the source
``memoryRecordId`` as the ``requestIdentifier`` so replays are idempotent.
A consumer (:mod:`stream_consumer`, run as a Lambda or locally) applies each
create, update, and delete to the target. A :mod:`record_map` tracks which
target record each source record became, and creates send a ``clientToken``
derived from the stream event, so replays don't duplicate records.
* **STM (short-term events) via dual-write ``CreateEvent``** —
:class:`DualRegionEventWriter` writes every conversation turn to both regions:
the source normally (triggering extraction, which feeds the stream above) and
Expand All @@ -20,22 +21,27 @@
See ``README.md`` for deployment instructions.
"""

from .dual_writer import DualRegionEventWriter
from .record_map import DynamoDBRecordMap, InMemoryRecordMap, MappedRecord, RecordMap
from .stream_consumer import (
StreamStats,
make_target_client,
process_kinesis_records,
replicate_stream_event,
stream_delivery_resources,
)
from .dual_writer import DualRegionEventWriter

__all__ = [
# STM via dual-write CreateEvent (extractionMode="SKIP" on target)
"DualRegionEventWriter",
# LTM via record streaming
"DynamoDBRecordMap",
"InMemoryRecordMap",
"MappedRecord",
"RecordMap",
"StreamStats",
"make_target_client",
"process_kinesis_records",
"replicate_stream_event",
"stream_delivery_resources",
# STM via dual-write CreateEvent (extractionMode="SKIP" on target)
"DualRegionEventWriter",
]
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,14 @@
get a stable, idempotent event identity in both regions, pass an explicit
``clientToken`` and the same ``eventTimestamp`` to both writes — don't rely on
the SDK auto-generating either. :class:`DualRegionEventWriter` does this for you.

To retry a turn, pass the same ``client_token`` and ``event_timestamp`` again (for
example, values derived from your own turn ID). If you let ``record_turn``
generate them, a retry after a timeout writes the turn a second time.
"""

import logging
import uuid
from typing import Optional

import boto3
from botocore.exceptions import ClientError
Expand Down Expand Up @@ -57,7 +60,7 @@ def __init__(
target_memory_id: str,
source_region: str,
target_region: str,
session: Optional[boto3.Session] = None,
session: boto3.Session | None = None,
):
self.source_memory_id = source_memory_id
self.target_memory_id = target_memory_id
Expand All @@ -71,8 +74,8 @@ def record_turn(
session_id: str,
role: str,
text: str,
event_timestamp: Optional[float] = None,
client_token: Optional[str] = None,
event_timestamp: float | None = None,
client_token: str | None = None,
) -> dict:
"""Record one conversation turn in BOTH regions.

Expand Down Expand Up @@ -135,14 +138,14 @@ def _create_event_idempotent(
extraction_mode=None,
):
"""CreateEvent that treats an idempotent "already exists" collision as OK."""
kwargs = dict(
memoryId=memory_id,
actorId=actor_id,
sessionId=session_id,
eventTimestamp=event_timestamp,
payload=payload,
clientToken=client_token,
)
kwargs = {
"memoryId": memory_id,
"actorId": actor_id,
"sessionId": session_id,
"eventTimestamp": event_timestamp,
"payload": payload,
"clientToken": client_token,
}
if extraction_mode:
kwargs["extractionMode"] = extraction_mode
try:
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
"""Source-to-target record ID map for LTM replication.

The target memory assigns its own ``memoryRecordId`` to every record that the
consumer creates, and a ``MemoryRecordDeleted`` stream event carries only the
source IDs. To apply updates and deletes, the consumer has to remember which
target record each source record became, and the time of the last change that
it applied (so a redelivered or out-of-order event is skipped, not re-applied).

Two implementations share one small interface:

* :class:`DynamoDBRecordMap` — the Lambda path, backed by the table in
``infra/streaming-stack.yaml``.
* :class:`InMemoryRecordMap` — the local demo and the unit tests.
"""

from dataclasses import dataclass
from decimal import Decimal
from typing import Protocol


@dataclass(frozen=True)
class MappedRecord:
"""What the consumer knows about one replicated source record."""

target_record_id: str | None # None for a tombstone of a never-replicated record
event_time: float # epoch seconds of the last change applied
deleted: bool = False


class RecordMap(Protocol):
def get(self, source_record_id: str) -> MappedRecord | None: ...

def put(self, source_record_id: str, record: MappedRecord) -> None: ...


class InMemoryRecordMap:
"""Process-local map: fine for a demo loop, not for Lambda (each container has its own)."""

def __init__(self):
self._records = {}

def get(self, source_record_id: str) -> MappedRecord | None:
return self._records.get(source_record_id)

def put(self, source_record_id: str, record: MappedRecord) -> None:
self._records[source_record_id] = record


class DynamoDBRecordMap:
"""Durable map in a DynamoDB table keyed by ``sourceRecordId`` (a boto3 ``Table`` resource)."""

def __init__(self, table):
self._table = table

def get(self, source_record_id: str) -> MappedRecord | None:
item = self._table.get_item(Key={"sourceRecordId": source_record_id}, ConsistentRead=True).get("Item")
if not item:
return None
return MappedRecord(
target_record_id=item.get("targetRecordId"),
event_time=float(item["eventTime"]),
deleted=bool(item.get("deleted", False)),
)

def put(self, source_record_id: str, record: MappedRecord) -> None:
item = {
"sourceRecordId": source_record_id,
"eventTime": Decimal(str(record.event_time)),
"deleted": record.deleted,
}
if record.target_record_id:
item["targetRecordId"] = record.target_record_id
self._table.put_item(Item=item)
Loading
Loading