Skip to content
12 changes: 12 additions & 0 deletions .github/workflows/rust.yml
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,16 @@ jobs:
- "chapter-08/lambdas/process_link_created"
- "chapter-08/lambdas/process_link_clicked"
- "chapter-08/shared"
- "chapter-09/examples/lambda-process-sqs"
- "chapter-09/examples/otel-simple"
- "chapter-09/examples/sqs-receive-traced-message"
- "chapter-09/examples/sqs-send-traced-message"
- "chapter-09/lambdas/create_link"
- "chapter-09/lambdas/get_links"
- "chapter-09/lambdas/visit_link"
- "chapter-09/lambdas/process_link_created"
- "chapter-09/lambdas/process_link_clicked"
- "chapter-09/shared"
fail-fast: false

steps:
Expand All @@ -62,4 +72,6 @@ jobs:
working-directory: ${{ matrix.project }}
- name: Run tests
run: cargo test --verbose
env:
ENV: test
working-directory: ${{ matrix.project }}
5 changes: 1 addition & 4 deletions chapter-04-challenge/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,7 @@ serde_json = "1.0"
cuid2 = "0.1"
aws-config = { version = "1.1", features = ["behavior-version-latest"] }
aws-sdk-dynamodb = "1.31"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
scraper = "0.16"


Expand Down
5 changes: 1 addition & 4 deletions chapter-04/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,7 @@ serde_json = "1.0"
cuid2 = "0.1"
aws-config = { version = "1.1", features = ["behavior-version-latest"] }
aws-sdk-dynamodb = "1.31"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
scraper = "0.22"

[dev-dependencies]
5 changes: 1 addition & 4 deletions chapter-05-challenge/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,5 @@ cuid2 = "0.1"
serde = "1.0"
serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.11"
5 changes: 1 addition & 4 deletions chapter-05/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,8 +9,5 @@ cuid2 = "0.1"
serde = "1.0"
serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.11"
5 changes: 1 addition & 4 deletions chapter-06-challenge/integration-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,4 @@ serde_json = "1.0"
shared = { path = "../shared" }
aws-config = { version = "1.1.7", features = ["behavior-version-latest"] }
aws-sdk-cloudformation = "1.41.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
5 changes: 1 addition & 4 deletions chapter-06-challenge/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,7 @@ cuid2 = "0.1"
serde = "1.0"
serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.11"
async-trait = "0.1.81"
mockall = { version = "0.13", optional = true }
Expand Down
5 changes: 1 addition & 4 deletions chapter-06/integration-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,4 @@ serde_json = "1.0"
shared = { path = "../shared" }
aws-config = { version = "1.1.7", features = ["behavior-version-latest"] }
aws-sdk-cloudformation = "1.41.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
5 changes: 1 addition & 4 deletions chapter-06/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -13,10 +13,7 @@ cuid2 = "0.1"
serde = "1.0"
serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.11"
async-trait = "0.1.81"
mockall = { version = "0.13", optional = true }
Expand Down
5 changes: 1 addition & 4 deletions chapter-07/01-sdks/integration-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,4 @@ serde_json = "1.0"
shared = { path = "../shared" }
aws-config = { version = "1.1.7", features = ["behavior-version-latest"] }
aws-sdk-cloudformation = "1.41.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
5 changes: 1 addition & 4 deletions chapter-07/01-sdks/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,7 @@ serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
aws-sdk-ssm = "1.31"
aws-sdk-secretsmanager = "1.66.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.14"
async-trait = "0.1.81"
mockall = { version = "0.13", optional = true }
Expand Down
5 changes: 1 addition & 4 deletions chapter-07/02-lambda-extension/integration-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,4 @@ serde_json = "1.0"
shared = { path = "../shared" }
aws-config = { version = "1.1.7", features = ["behavior-version-latest"] }
aws-sdk-cloudformation = "1.41.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
5 changes: 1 addition & 4 deletions chapter-07/02-lambda-extension/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,7 @@ serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
aws-sdk-ssm = "1.31"
aws-sdk-secretsmanager = "1.66.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.14"
async-trait = "0.1.81"
mockall = { version = "0.13", optional = true }
Expand Down
5 changes: 1 addition & 4 deletions chapter-08/integration-tests/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,4 @@ serde_json = "1.0"
shared = { path = "../shared" }
aws-config = { version = "1.1.7", features = ["behavior-version-latest"] }
aws-sdk-cloudformation = "1.41.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
5 changes: 1 addition & 4 deletions chapter-08/shared/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,10 +15,7 @@ serde_json = "1.0"
aws-sdk-dynamodb = "1.31"
aws-sdk-ssm = "1.31"
aws-sdk-secretsmanager = "1.66.0"
reqwest = { version = "0.12", default-features = false, features = [
"rustls-tls",
"http2",
] }
reqwest = "0.13"
lambda_http = "0.14"
async-trait = "0.1.81"
mockall = { version = "0.13", optional = true }
Expand Down
4 changes: 3 additions & 1 deletion chapter-09/Cargo.toml
Original file line number Diff line number Diff line change
@@ -1,7 +1,9 @@
[workspace]
resolver = "2"
members = [
"examples/sqs-sendmessage-structured",
"examples/otel-simple",
"examples/sqs-receive-traced-message",
"examples/sqs-send-traced-message",
"examples/lambda-process-sqs",
"shared",
"lambdas/create_link",
Expand Down
15 changes: 15 additions & 0 deletions chapter-09/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,15 @@
# Query examples for Amazon CloudWatch

## OTEL in CloudWatch

1. Enable [CloudWatch Transaction Search](https://docs.aws.amazon.com/AmazonCloudWatch/latest/monitoring/Enable-TransactionSearch.html). For testing purposes, set sampling to 100%.
2. View all process spans by querying `name ^ process`

## CloudWatch Log Insights

```
fields @timestamp, resource.attributes.service.name, severityText, body
| filter resource.attributes.faas.name = "GetLinksFunction-dev"
| sort @timestamp desc
| limit 10000
```
11 changes: 11 additions & 0 deletions chapter-09/docker-compose.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
services:
jaeger:
image: jaegertracing/opentelemetry-all-in-one:latest
ports:
- 16686:16686
- 13133:13133
- 4317:4317
- 4318:4318
network_mode: "host"
security_opt:
- no-new-privileges:true
6 changes: 6 additions & 0 deletions chapter-09/examples/lambda-process-sqs/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,10 +4,16 @@ version = "0.1.0"
edition = "2021"

[dependencies]
shared = { path = "../../shared" }

aws_lambda_events = { version = "1.0.0", default-features = false, features = [
"sqs",
] }
serde_json = "1.0"
serde = "1.0.228"
lambda_runtime = "1.0.1"
tokio = { version = "1", features = ["macros"] }
cloudevents-sdk = "0.9.0"

opentelemetry = "0.31.0"
tracing = "0.1.43"
47 changes: 35 additions & 12 deletions chapter-09/examples/lambda-process-sqs/src/event_handler.rs
Original file line number Diff line number Diff line change
@@ -1,30 +1,53 @@
use aws_lambda_events::event::sqs::SqsEvent;
use lambda_runtime::{Error, LambdaEvent};
use ::tracing::Span;
use aws_lambda_events::{event::sqs::SqsEvent, sqs::SqsMessage};
use lambda_runtime::{tracing, Error, LambdaEvent};
use serde::Deserialize;
use shared::observability::add_span_link_from;

#[derive(Deserialize)]
struct MyMessage {
task_id: String,
data: String,
}

#[tracing::instrument(skip(event))]
pub async fn function_handler(event: LambdaEvent<SqsEvent>) -> Result<(), Error> {
for record in event.payload.records {
// Get the message body
let body = record.body.unwrap_or_default();

// Parse the JSON message
let message: MyMessage = serde_json::from_str(&body)?;

// Process the message
println!("Processing task: {} - {}", message.task_id, message.data);
process_task(message).await?;
process_message(record).await?;
}

Ok(())
}

async fn process_task(message: MyMessage) -> Result<(), Error> {
#[tracing::instrument(skip(message))]
async fn process_message(message: SqsMessage) -> Result<(), Error> {
// Get the message body
let current_span = Span::current();
tracing::info!("Received message: {:?}", message.body);

let cloud_event: cloudevents::Event =
match serde_json::from_str(message.body.as_ref().unwrap_or(&"".to_string())) {
Ok(event) => event,
Err(e) => {
tracing::error!("Failed to deserialize CloudEvent: {:?}", e);
return Err(Error::from(e));
}
};

add_span_link_from(&current_span, &cloud_event);

// Parse the JSON message
let cloud_event_data = cloud_event.data().ok_or("CloudEvent has no data")?;

let message: MyMessage = match cloud_event_data {
cloudevents::Data::Binary(items) => serde_json::from_slice(items)?,
cloudevents::Data::String(string_data) => serde_json::from_str(&string_data)?,
cloudevents::Data::Json(value) => serde_json::from_value(value.clone())?,
};

// Process the message
println!("Processing task: {} - {}", message.task_id, message.data);

// Your business logic here
println!("Task {} completed", message.task_id);
Ok(())
Expand Down
16 changes: 13 additions & 3 deletions chapter-09/examples/lambda-process-sqs/src/main.rs
Original file line number Diff line number Diff line change
@@ -1,11 +1,21 @@
use std::sync::Arc;

use event_handler::function_handler;
use lambda_runtime::{run, service_fn, tracing, Error};
use lambda_runtime::{run, service_fn, Error};

mod event_handler;

#[tokio::main]
async fn main() -> Result<(), Error> {
tracing::init_default_subscriber();
let otel_guard =
Arc::new(shared::observability::init_otel().expect("Failed to initialize telemetry"));

run(service_fn(|evt| async {
let res = function_handler(evt).await;

otel_guard.flush();

run(service_fn(function_handler)).await
res
}))
.await
}
13 changes: 13 additions & 0 deletions chapter-09/examples/otel-simple/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
[package]
name = "otel-simple"
version = "0.1.0"
edition = "2021"

[dependencies]
shared = { path = "../../shared" }
tokio = { version = "1.38", features = ["macros", "rt-multi-thread"] }
serde_json = "1.0"
serde = "1.0.228"

opentelemetry = "0.31.0"
tracing = "0.1.43"
29 changes: 29 additions & 0 deletions chapter-09/examples/otel-simple/src/main.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
use shared::observability::init_otel;
use tracing::{error, info, warn};

#[tokio::main]
async fn main() {
// Telemetry is flushed on drop
let _otel_guard = init_otel().expect("Failed to initialize telemetry");

do_work();
}

#[tracing::instrument()]
fn do_work() {
// Simulate some work being done
std::thread::sleep(std::time::Duration::from_millis(500));

let my_random_variable = "this is the value";

error!("the random variable value is {}.", my_random_variable);
warn!("This is a warning log");
info!("Work completed");

do_some_more_work();
}

#[tracing::instrument()]
fn do_some_more_work() {
std::thread::sleep(std::time::Duration::from_millis(1500));
}
17 changes: 17 additions & 0 deletions chapter-09/examples/sqs-receive-traced-message/Cargo.toml
Original file line number Diff line number Diff line change
@@ -0,0 +1,17 @@
[package]
name = "sqs-receive-traced-message"
version = "0.1.0"
edition = "2021"

[dependencies]
shared = { path = "../../shared" }
tokio = { version = "1.38", features = ["macros", "rt-multi-thread"] }
aws-config = { version = "1.1", features = ["behavior-version-latest"] }
aws-sdk-sqs = "1.90.0"
serde_json = "1.0"
serde = "1.0.228"
cloudevents-sdk = "0.9.0"
cuid2 = "0.1"

opentelemetry = "0.31.0"
tracing = "0.1.43"
Loading
Loading