Skip to content
Merged
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
1 change: 1 addition & 0 deletions src/sink/memory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,7 @@ mod tests {
payload: serde_json::json!({ "test": id }),
metadata: None,
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
}
}
Expand Down
2 changes: 2 additions & 0 deletions src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,7 @@ impl<S: SchemaStore> StreamStore<S> {
payload,
metadata: row.event_metadata,
stream_id: crate::types::StreamId::from(event_stream_id as u64),
commit_lsn: None,
lsn: row.event_lsn.and_then(|s| s.parse().ok()),
}));
}
Expand Down Expand Up @@ -176,6 +177,7 @@ impl<S: SchemaStore> StreamStore<S> {
payload: row.payload,
metadata: row.metadata,
stream_id: crate::types::StreamId::from(row.stream_id as u64),
commit_lsn: None,
lsn: row.lsn.and_then(|s| s.parse().ok()),
});

Expand Down
54 changes: 36 additions & 18 deletions src/types/event.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,8 +49,12 @@ pub struct TriggeredEvent {
pub payload: serde_json::Value,
pub metadata: Option<serde_json::Value>,
pub stream_id: StreamId,
/// The logical replication transaction commit LSN.
/// Available for live events and safe to use as a replica consistency watermark.
pub commit_lsn: Option<PgLsn>,
/// The WAL LSN at the time the event was inserted.
/// Used for precise replay point lookup during slot recovery.
/// Used for precise replay point lookup during slot recovery, but not safe as a
/// replica consistency watermark because it precedes the transaction commit.
pub lsn: Option<PgLsn>,
}

Expand Down Expand Up @@ -78,6 +82,7 @@ macro_rules! missing {
pub fn convert_event_from_table(
table_row: &mut TableRow,
column_schemas: &[ColumnSchema],
commit_lsn: Option<PgLsn>,
) -> EtlResult<TriggeredEvent> {
let mut id = None;
let mut created_at = None;
Expand Down Expand Up @@ -110,6 +115,7 @@ pub fn convert_event_from_table(
payload: payload.ok_or_else(|| missing!("payload"))?,
stream_id: stream_id.ok_or_else(|| missing!("stream_id"))?,
metadata,
commit_lsn,
lsn,
})
}
Expand All @@ -120,7 +126,7 @@ pub fn convert_events_from_table_rows(
) -> EtlResult<Vec<TriggeredEvent>> {
table_rows
.into_iter()
.map(|mut table_row| convert_event_from_table(&mut table_row, column_schemas))
.map(|mut table_row| convert_event_from_table(&mut table_row, column_schemas, None))
.collect()
}

Expand All @@ -131,10 +137,15 @@ pub fn convert_stream_events_from_events(
events
.into_iter()
.filter_map(|event| match event {
Event::Insert(mut insert_event) => Some(convert_event_from_table(
&mut insert_event.table_row,
column_schemas,
)),
Event::Insert(mut insert_event) => {
let commit_lsn = Some(insert_event.commit_lsn);

Some(convert_event_from_table(
&mut insert_event.table_row,
column_schemas,
commit_lsn,
))
}
Event::Begin(_)
| Event::Commit(_)
| Event::Update(_)
Expand Down Expand Up @@ -175,9 +186,9 @@ mod tests {
Cell::Uuid(id),
Cell::TimestampTz(created_at),
Cell::Json(payload),
Cell::Null, // metadata
Cell::I64(1), // stream_id
Cell::String("0/16B3748".to_string()), // lsn (parsed to PgLsn)
Cell::Null, // metadata
Cell::I64(1), // stream_id
Cell::String("0/100".to_string()), // lsn (parsed to PgLsn)
],
}
}
Expand All @@ -200,20 +211,21 @@ mod tests {
}

#[test]
fn test_convert_event_from_table_valid() {
fn test_convert_replayed_event_preserves_row_lsn_without_commit_lsn() {
let column_schemas = make_column_schemas();
let id = Uuid::new_v4();
let created_at = Utc::now();
let payload = serde_json::json!({"test": "data"});

let mut table_row = make_table_row(id, created_at, payload.clone());

let result = convert_event_from_table(&mut table_row, &column_schemas).unwrap();
let result = convert_event_from_table(&mut table_row, &column_schemas, None).unwrap();

assert_eq!(result.id.id, id.to_string());
assert_eq!(result.id.created_at, created_at);
assert_eq!(result.payload, payload);
assert_eq!(result.lsn, Some("0/16B3748".parse().unwrap()));
assert_eq!(result.lsn, Some("0/100".parse().unwrap()));
assert_eq!(result.commit_lsn, None);
}

#[test]
Expand All @@ -225,12 +237,13 @@ mod tests {

let mut table_row = make_table_row_without_lsn(id, created_at, payload.clone());

let result = convert_event_from_table(&mut table_row, &column_schemas).unwrap();
let result = convert_event_from_table(&mut table_row, &column_schemas, None).unwrap();

assert_eq!(result.id.id, id.to_string());
assert_eq!(result.id.created_at, created_at);
assert_eq!(result.payload, payload);
assert_eq!(result.lsn, None);
assert_eq!(result.commit_lsn, None);
}

#[test]
Expand All @@ -247,7 +260,7 @@ mod tests {
],
};

let result = convert_event_from_table(&mut table_row, &column_schemas);
let result = convert_event_from_table(&mut table_row, &column_schemas, None);

assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("Missing id"));
Expand All @@ -267,7 +280,7 @@ mod tests {
],
};

let result = convert_event_from_table(&mut table_row, &column_schemas);
let result = convert_event_from_table(&mut table_row, &column_schemas, None);

assert!(result.is_err());
assert!(
Expand All @@ -292,7 +305,7 @@ mod tests {
],
};

let result = convert_event_from_table(&mut table_row, &column_schemas);
let result = convert_event_from_table(&mut table_row, &column_schemas, None);

assert!(result.is_err());
assert!(result.unwrap_err().to_string().contains("Missing payload"));
Expand Down Expand Up @@ -330,15 +343,15 @@ mod tests {
}

#[test]
fn test_convert_stream_events_from_events_filters_inserts_only() {
fn test_convert_live_insert_preserves_commit_lsn() {
let column_schemas = make_column_schemas();
let id = Uuid::new_v4();
let ts = Utc::now();

let events = vec![
Event::Insert(InsertEvent {
start_lsn: PgLsn::from(0),
commit_lsn: PgLsn::from(0),
commit_lsn: "0/200".parse().unwrap(),
table_id: TableId::new(1),
table_row: make_table_row(id, ts, serde_json::json!({"test": 1})),
}),
Expand All @@ -353,6 +366,8 @@ mod tests {
let first = result.first().expect("first element exists");
assert_eq!(first.id.id, id.to_string());
assert_eq!(first.payload.get("test").and_then(|v| v.as_i64()), Some(1));
assert_eq!(first.lsn, Some("0/100".parse().unwrap()));
assert_eq!(first.commit_lsn, Some("0/200".parse().unwrap()));
}

#[test]
Expand Down Expand Up @@ -455,13 +470,15 @@ mod tests {
payload: payload.clone(),
metadata: None,
stream_id: StreamId::from(1u64),
commit_lsn: Some("0/200".parse().unwrap()),
lsn: Some("0/16B3748".parse().unwrap()),
};
let event2 = TriggeredEvent {
id,
payload,
metadata: None,
stream_id: StreamId::from(1u64),
commit_lsn: Some("0/200".parse().unwrap()),
lsn: Some("0/16B3748".parse().unwrap()),
};

Expand All @@ -479,6 +496,7 @@ mod tests {
payload,
metadata: None,
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: None,
};
let (id_returned, created_at) = event.primary_keys();
Expand Down
3 changes: 3 additions & 0 deletions tests/elasticsearch_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ fn make_test_event(key: &str) -> TriggeredEvent {
stream_id: StreamId::default(),
payload: serde_json::json!({ "key": key, "value": "test_data" }),
metadata: Some(serde_json::json!({ "source": "test" })),
commit_lsn: None,
lsn: Some(PgLsn::from(12345u64)),
}
}
Expand Down Expand Up @@ -131,6 +132,7 @@ async fn test_elasticsearch_sink_indexes_only_payload() {
stream_id: StreamId::default(),
payload: serde_json::json!({ "action": "created", "user_id": 456 }),
metadata: Some(serde_json::json!({ "source": "api" })),
commit_lsn: None,
lsn: Some(PgLsn::from(99999u64)),
};
let event_id = event.id.id.clone();
Expand Down Expand Up @@ -247,6 +249,7 @@ async fn test_elasticsearch_sink_uses_index_from_metadata() {
stream_id: StreamId::default(),
payload: serde_json::json!({ "routed": true }),
metadata: Some(serde_json::json!({ "index": metadata_index })),
commit_lsn: None,
lsn: None,
};
let event_id = event.id.id.clone();
Expand Down
2 changes: 2 additions & 0 deletions tests/gcp_pubsub_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ fn make_test_event(key: &str) -> TriggeredEvent {
stream_id: StreamId::default(),
payload: serde_json::json!({ "key": key, "value": "test_data" }),
metadata: Some(serde_json::json!({ "source": "test" })),
commit_lsn: None,
lsn: Some(PgLsn::from(12345u64)),
}
}
Expand Down Expand Up @@ -201,6 +202,7 @@ async fn test_gcp_pubsub_sink_sends_only_payload() {
stream_id: StreamId::default(),
payload: serde_json::json!({ "action": "created", "user_id": 123 }),
metadata: Some(serde_json::json!({ "routing_key": "orders" })),
commit_lsn: None,
lsn: Some(PgLsn::from(99999u64)),
};

Expand Down
2 changes: 2 additions & 0 deletions tests/kafka_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ fn make_test_event(id: &str) -> TriggeredEvent {
}),
metadata: Some(serde_json::json!({ "source": "test" })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
}
}
Expand Down Expand Up @@ -148,6 +149,7 @@ async fn test_kafka_sink_uses_topic_from_metadata() {
}),
metadata: Some(serde_json::json!({ "topic": topic })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
};

Expand Down
2 changes: 2 additions & 0 deletions tests/kinesis_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ fn make_test_event(key: &str) -> TriggeredEvent {
stream_id: StreamId::default(),
payload: serde_json::json!({ "key": key, "value": "test_data" }),
metadata: Some(serde_json::json!({ "source": "test" })),
commit_lsn: None,
lsn: Some(PgLsn::from(12345u64)),
}
}
Expand Down Expand Up @@ -231,6 +232,7 @@ async fn test_kinesis_sink_uses_stream_from_metadata() {
stream_id: StreamId::default(),
payload: serde_json::json!({ "action": "created" }),
metadata: Some(serde_json::json!({ "stream": stream_name })),
commit_lsn: None,
lsn: Some(PgLsn::from(99999u64)),
};

Expand Down
3 changes: 3 additions & 0 deletions tests/meilisearch_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ fn make_test_event(key: &str) -> TriggeredEvent {
stream_id: StreamId::default(),
payload: serde_json::json!({ "id": id, "key": key, "value": "test_data" }),
metadata: Some(serde_json::json!({ "source": "test" })),
commit_lsn: None,
lsn: Some(PgLsn::from(12345u64)),
}
}
Expand Down Expand Up @@ -119,6 +120,7 @@ async fn test_meilisearch_sink_indexes_only_payload() {
stream_id: StreamId::default(),
payload: serde_json::json!({ "id": event_id, "action": "created", "user": 456 }),
metadata: Some(serde_json::json!({ "source": "api" })),
commit_lsn: None,
lsn: Some(PgLsn::from(99999u64)),
};

Expand Down Expand Up @@ -217,6 +219,7 @@ async fn test_meilisearch_sink_uses_index_from_metadata() {
stream_id: StreamId::default(),
payload: serde_json::json!({ "id": event_id, "routed": true }),
metadata: Some(serde_json::json!({ "index": metadata_index })),
commit_lsn: None,
lsn: None,
};

Expand Down
2 changes: 2 additions & 0 deletions tests/nats_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ fn make_test_event(id: &str) -> TriggeredEvent {
}),
metadata: Some(serde_json::json!({ "source": "test" })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
}
}
Expand Down Expand Up @@ -148,6 +149,7 @@ async fn test_nats_sink_uses_topic_from_metadata() {
}),
metadata: Some(serde_json::json!({ "topic": subject })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
};

Expand Down
2 changes: 2 additions & 0 deletions tests/rabbitmq_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ fn make_test_event(id: &str) -> TriggeredEvent {
}),
metadata: Some(serde_json::json!({ "source": "test" })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
}
}
Expand Down Expand Up @@ -195,6 +196,7 @@ async fn test_rabbitmq_sink_uses_metadata_routing() {
"routing_key": routing_key
})),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
};

Expand Down
2 changes: 2 additions & 0 deletions tests/redis_streams_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ fn make_test_event(id: &str) -> TriggeredEvent {
}),
metadata: Some(serde_json::json!({ "source": "test" })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
}
}
Expand Down Expand Up @@ -172,6 +173,7 @@ async fn test_redis_streams_sink_uses_stream_from_metadata() {
}),
metadata: Some(serde_json::json!({ "stream": stream_name })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
};

Expand Down
2 changes: 2 additions & 0 deletions tests/redis_strings_sink_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ fn make_test_event(id: &str) -> TriggeredEvent {
}),
metadata: None,
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
}
}
Expand Down Expand Up @@ -149,6 +150,7 @@ async fn test_redis_strings_sink_uses_key_from_metadata() {
}),
metadata: Some(serde_json::json!({ "key": custom_key })),
stream_id: StreamId::from(1u64),
commit_lsn: None,
lsn: Some("0/16B3748".parse().unwrap()),
};

Expand Down
Loading
Loading