diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-common/src/main/java/datadog/trace/instrumentation/aws/AwsAccountIdentity.java b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/main/java/datadog/trace/instrumentation/aws/AwsAccountIdentity.java
new file mode 100644
index 00000000000..56bfc88d218
--- /dev/null
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/main/java/datadog/trace/instrumentation/aws/AwsAccountIdentity.java
@@ -0,0 +1,182 @@
+package datadog.trace.instrumentation.aws;
+
+/**
+ * Derives AWS account identity for resources addressed by AWS SDK requests.
+ *
+ *
Two sources are trustworthy and used here:
+ *
+ *
+ * An ARN carried by the request itself (for example a DynamoDB {@code TableName} given as an
+ * ARN, or an SNS {@code TopicArn}), parsed once with {@link AwsArn}.
+ * The account that owns the credentials signing the request. Several AWS services resolve a
+ * bare resource name in the requestor's own account (DynamoDB documents this explicitly: "If
+ * you only provide the table name parameter instead of a complete ARN, the API operation will
+ * be performed on the table in the account to which the requestor belongs"), so for those
+ * services the caller account is the resource owner.
+ *
+ *
+ * The caller account is taken from the credentials object when the SDK exposes it ({@code
+ * AwsCredentialsIdentity.accountId()} in AWS SDK for Java v2 2.26+, populated by the STS, SSO,
+ * profile, process and container credential providers). As an opt-in fallback for older SDKs the
+ * account can be decoded from the access key ID, which encodes it in its trailing characters.
+ */
+public final class AwsAccountIdentity {
+
+ private AwsAccountIdentity() {}
+
+ /** {@code true} for exactly twelve ASCII digits. */
+ public static boolean isAccountId(final String value) {
+ if (value == null || value.length() != 12) {
+ return false;
+ }
+ for (int i = 0; i < 12; i++) {
+ char c = value.charAt(i);
+ if (c < '0' || c > '9') {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ /** Maps a Region code to its partition. Unknown prefixes map to the commercial partition. */
+ public static String partitionForRegion(final String region) {
+ if (region == null) {
+ return "aws";
+ }
+ if (region.startsWith("cn-")) {
+ return "aws-cn";
+ }
+ if (region.startsWith("us-gov-")) {
+ return "aws-us-gov";
+ }
+ if (region.startsWith("us-isob-")) {
+ return "aws-iso-b";
+ }
+ if (region.startsWith("us-isof-")) {
+ return "aws-iso-f";
+ }
+ if (region.startsWith("us-iso-")) {
+ return "aws-iso";
+ }
+ if (region.startsWith("eu-isoe-")) {
+ return "aws-iso-e";
+ }
+ if (region.startsWith("eusc-")) {
+ return "aws-eusc";
+ }
+ return "aws";
+ }
+
+ /** Builds a DynamoDB table ARN, or returns {@code null} when the Region or account is unknown. */
+ public static String dynamoDbTableArn(
+ final String region, final String account, final String tableName) {
+ if (region == null || !isAccountId(account) || tableName == null || tableName.isEmpty()) {
+ return null;
+ }
+ return "arn:"
+ + partitionForRegion(region)
+ + ":dynamodb:"
+ + region
+ + ':'
+ + account
+ + ":table/"
+ + tableName;
+ }
+
+ /**
+ * Decodes the owning account from an AWS access key ID.
+ *
+ *
Access key IDs are a 4 character prefix ({@code AKIA} for long-term keys, {@code ASIA} for
+ * temporary ones) followed by 16 base32 characters that encode 10 bytes. For keys issued since
+ * late March 2019 the account ID is held in the first 48 bits of those bytes, shifted left by 7,
+ * and the top bit of those 48 is always set (the fifth character of the key is {@code Q} or later
+ * in the base32 alphabet). Older keys carry no account at all and that bit is clear, so they are
+ * rejected rather than decoded into an arbitrary value. This encoding is not part of the
+ * documented AWS API surface, which is why callers only use it when explicitly enabled.
+ *
+ * @return the 12-digit account, or {@code null} when the key does not have the expected shape or
+ * predates the account-encoding format.
+ */
+ public static String accountFromAccessKeyId(final String accessKeyId) {
+ if (accessKeyId == null || accessKeyId.length() != 20) {
+ return null;
+ }
+ if (!(accessKeyId.startsWith("AKIA") || accessKeyId.startsWith("ASIA"))) {
+ return null;
+ }
+ byte[] decoded = decodeBase32(accessKeyId, 4, 20);
+ return decoded == null ? null : accountFromEncodedBytes(decoded);
+ }
+
+ /**
+ * Decodes RFC 4648 base32 (no padding, case-insensitive) from {@code value[from, to)}.
+ *
+ * @return the decoded bytes, or {@code null} when a character is outside the alphabet.
+ */
+ static byte[] decodeBase32(final CharSequence value, final int from, final int to) {
+ byte[] out = new byte[(to - from) * 5 / 8];
+ int buffer = 0;
+ int bits = 0;
+ int index = 0;
+ for (int i = from; i < to; i++) {
+ int digit = base32Value(value.charAt(i));
+ if (digit < 0) {
+ return null;
+ }
+ buffer = (buffer << 5) | digit;
+ bits += 5;
+ if (bits >= 8) {
+ bits -= 8;
+ out[index++] = (byte) ((buffer >> bits) & 0xff);
+ }
+ }
+ return out;
+ }
+
+ /** Set on the first 48 bits of every access key issued in the account-encoding format. */
+ static final long ACCOUNT_FORMAT_MARKER = 0x800000000000L;
+
+ /**
+ * Extracts the account ID from the decoded bytes of an access key ID: the first 6 bytes hold
+ * {@code account << 7} with the top bit set as the format marker.
+ *
+ * @return the zero-padded 12-digit account, or {@code null} when fewer than 6 bytes are given,
+ * the format marker is clear (a legacy key with no encoded account), or the value does not
+ * fit in 12 digits.
+ */
+ static String accountFromEncodedBytes(final byte[] decoded) {
+ if (decoded.length < 6) {
+ return null;
+ }
+ long firstSixBytes = 0;
+ for (int i = 0; i < 6; i++) {
+ firstSixBytes = (firstSixBytes << 8) | (decoded[i] & 0xff);
+ }
+ if ((firstSixBytes & ACCOUNT_FORMAT_MARKER) == 0) {
+ return null;
+ }
+ long account = (firstSixBytes & 0x7fffffffff80L) >>> 7;
+ String text = Long.toString(account);
+ if (text.length() > 12) {
+ return null;
+ }
+ StringBuilder padded = new StringBuilder(12);
+ for (int i = text.length(); i < 12; i++) {
+ padded.append('0');
+ }
+ return padded.append(text).toString();
+ }
+
+ private static int base32Value(final char c) {
+ if (c >= 'A' && c <= 'Z') {
+ return c - 'A';
+ }
+ if (c >= 'a' && c <= 'z') {
+ return c - 'a';
+ }
+ if (c >= '2' && c <= '7') {
+ return c - '2' + 26;
+ }
+ return -1;
+ }
+}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-common/src/main/java/datadog/trace/instrumentation/aws/AwsArn.java b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/main/java/datadog/trace/instrumentation/aws/AwsArn.java
new file mode 100644
index 00000000000..05b0f985aa2
--- /dev/null
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/main/java/datadog/trace/instrumentation/aws/AwsArn.java
@@ -0,0 +1,125 @@
+package datadog.trace.instrumentation.aws;
+
+/**
+ * A parsed Amazon Resource Name: {@code arn:::::}.
+ *
+ * Parsed once so callers that need several fields do not re-scan the string. The AWS SDK ships
+ * its own parser ({@code software.amazon.awssdk.arns.Arn}) but it lives in a module that neither
+ * the SDK v2 2.2.0 floor these instrumentations compile against nor SDK v1 provide, so a minimal
+ * one is kept here.
+ */
+public final class AwsArn {
+
+ private static final String TABLE_PREFIX = "table/";
+
+ private final String raw;
+ private final String partition;
+ private final String service;
+ private final String region;
+ private final String account;
+ private final String resource;
+
+ private AwsArn(
+ final String raw,
+ final String partition,
+ final String service,
+ final String region,
+ final String account,
+ final String resource) {
+ this.raw = raw;
+ this.partition = partition;
+ this.service = service;
+ this.region = region;
+ this.account = account;
+ this.resource = resource;
+ }
+
+ /**
+ * Parses an ARN.
+ *
+ * @return the parsed ARN, or {@code null} when the value is not of the form {@code
+ * arn:partition:service:region:account:resource} (partition and service non-empty).
+ */
+ public static AwsArn parse(final String value) {
+ if (value == null || !value.startsWith("arn:")) {
+ return null;
+ }
+ int c1 = value.indexOf(':', 4);
+ if (c1 < 0) {
+ return null;
+ }
+ int c2 = value.indexOf(':', c1 + 1);
+ if (c2 < 0) {
+ return null;
+ }
+ int c3 = value.indexOf(':', c2 + 1);
+ if (c3 < 0) {
+ return null;
+ }
+ int c4 = value.indexOf(':', c3 + 1);
+ if (c4 < 0) {
+ return null;
+ }
+ if (c1 == 4 || c2 == c1 + 1 || c4 == value.length() - 1) {
+ // empty partition, empty service or empty resource
+ return null;
+ }
+ return new AwsArn(
+ value,
+ value.substring(4, c1),
+ value.substring(c1 + 1, c2),
+ emptyToNull(value.substring(c2 + 1, c3)),
+ emptyToNull(value.substring(c3 + 1, c4)),
+ value.substring(c4 + 1));
+ }
+
+ /** The original string. */
+ public String raw() {
+ return raw;
+ }
+
+ /** {@code aws}, {@code aws-cn}, {@code aws-us-gov}, ... */
+ public String partition() {
+ return partition;
+ }
+
+ /** {@code dynamodb}, {@code sns}, {@code s3}, ... */
+ public String service() {
+ return service;
+ }
+
+ /** The Region, or {@code null} for global services (S3, IAM). */
+ public String region() {
+ return region;
+ }
+
+ /** The 12-digit account, or {@code null} when absent (S3 ARNs) or not exactly twelve digits. */
+ public String account() {
+ return AwsAccountIdentity.isAccountId(account) ? account : null;
+ }
+
+ /** Everything after the fifth colon. Never empty. */
+ public String resource() {
+ return resource;
+ }
+
+ /**
+ * The table name of a DynamoDB table ARN, with any sub-resource ({@code /index/}, {@code
+ * /stream/}, ...) removed.
+ *
+ * @return the bare table name, or {@code null} when the resource is not a table.
+ */
+ public String dynamoDbTableName() {
+ if (!resource.startsWith(TABLE_PREFIX)) {
+ return null;
+ }
+ int start = TABLE_PREFIX.length();
+ int slash = resource.indexOf('/', start);
+ String name = slash < 0 ? resource.substring(start) : resource.substring(start, slash);
+ return name.isEmpty() ? null : name;
+ }
+
+ private static String emptyToNull(final String value) {
+ return value.isEmpty() ? null : value;
+ }
+}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-common/src/test/java/datadog/trace/instrumentation/aws/AwsAccountIdentityTest.java b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/test/java/datadog/trace/instrumentation/aws/AwsAccountIdentityTest.java
new file mode 100644
index 00000000000..8a011bb0530
--- /dev/null
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/test/java/datadog/trace/instrumentation/aws/AwsAccountIdentityTest.java
@@ -0,0 +1,183 @@
+package datadog.trace.instrumentation.aws;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import java.nio.charset.StandardCharsets;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.NullAndEmptySource;
+import org.junit.jupiter.params.provider.ValueSource;
+
+class AwsAccountIdentityTest {
+
+ @ParameterizedTest(name = "{0} is an account id")
+ @ValueSource(strings = {"123456789012", "000000000000", "999999999999"})
+ void acceptsTwelveDigits(String value) {
+ assertTrue(AwsAccountIdentity.isAccountId(value));
+ }
+
+ @ParameterizedTest(name = "{0} is not an account id")
+ @NullAndEmptySource
+ @ValueSource(
+ strings = {"12345", "1234567890123", "12345678901a", " 23456789012", "arn:aws:s3:::b"})
+ void rejectsNonAccountIds(String value) {
+ assertFalse(AwsAccountIdentity.isAccountId(value));
+ }
+
+ @ParameterizedTest(name = "{0} is in partition {1}")
+ @CsvSource(
+ nullValues = "NULL",
+ value = {
+ "us-east-1, aws",
+ "eu-central-1, aws",
+ "cn-north-1, aws-cn",
+ "us-gov-west-1, aws-us-gov",
+ "us-iso-east-1, aws-iso",
+ "us-isob-east-1, aws-iso-b",
+ "eu-isoe-west-1, aws-iso-e",
+ "us-isof-south-1, aws-iso-f",
+ "eusc-de-east-1, aws-eusc",
+ "NULL, aws",
+ })
+ void mapsRegionToPartition(String region, String partition) {
+ assertEquals(partition, AwsAccountIdentity.partitionForRegion(region));
+ }
+
+ @Test
+ void buildsDynamoDbTableArn() {
+ assertEquals(
+ "arn:aws:dynamodb:us-west-2:123456789012:table/orders",
+ AwsAccountIdentity.dynamoDbTableArn("us-west-2", "123456789012", "orders"));
+ assertEquals(
+ "arn:aws-cn:dynamodb:cn-north-1:123456789012:table/orders",
+ AwsAccountIdentity.dynamoDbTableArn("cn-north-1", "123456789012", "orders"));
+ assertNull(AwsAccountIdentity.dynamoDbTableArn(null, "123456789012", "orders"));
+ assertNull(AwsAccountIdentity.dynamoDbTableArn("us-west-2", null, "orders"));
+ assertNull(AwsAccountIdentity.dynamoDbTableArn("us-west-2", "1234", "orders"));
+ assertNull(AwsAccountIdentity.dynamoDbTableArn("us-west-2", "123456789012", ""));
+ }
+
+ @Test
+ void decodesAccountFromAccessKeyId() {
+ // Key IDs built by encoding known accounts with the documented layout: 4 char prefix, then
+ // base32 of 10 bytes whose first 48 bits are (account << 7) plus arbitrary low bits.
+ assertEquals(
+ "123456789012",
+ AwsAccountIdentity.accountFromAccessKeyId(syntheticKey("AKIA", 123456789012L)));
+ assertEquals(
+ "640168429175",
+ AwsAccountIdentity.accountFromAccessKeyId(syntheticKey("ASIA", 640168429175L)));
+ assertEquals(
+ "000000000001", AwsAccountIdentity.accountFromAccessKeyId(syntheticKey("ASIA", 1L)));
+ assertEquals(
+ "999999999999",
+ AwsAccountIdentity.accountFromAccessKeyId(syntheticKey("AKIA", 999999999999L)));
+ }
+
+ @ParameterizedTest(name = "rejects access key id {0}")
+ @NullAndEmptySource
+ @ValueSource(
+ strings = {
+ "my-access-key",
+ "AKIA",
+ // 19 chars
+ "AKIAIOSFODNN7EXAMPL",
+ // 21 chars
+ "AKIAIOSFODNN7EXAMPLEX",
+ // unknown prefix
+ "ABCDIOSFODNN7EXAMPLE",
+ // '1' is not base32
+ "AKIAIOSFODNN7EXAMPL1",
+ // '0' is not base32 (and wrong case prefix)
+ "akiaiosfodnn7example0",
+ })
+ void rejectsMalformedAccessKeyIds(String value) {
+ assertNull(AwsAccountIdentity.accountFromAccessKeyId(value));
+ }
+
+ @Test
+ void decodesRfc4648Base32() {
+ // RFC 4648 test vector: "foobar" -> MZXW6YTBOI (unpadded)
+ assertArrayEquals(
+ "foobar".getBytes(StandardCharsets.US_ASCII),
+ AwsAccountIdentity.decodeBase32("MZXW6YTBOI", 0, 10));
+ // lower case alphabet is accepted
+ assertArrayEquals(
+ "foobar".getBytes(StandardCharsets.US_ASCII),
+ AwsAccountIdentity.decodeBase32("mzxw6ytboi", 0, 10));
+ // sub-range decoding, as used to skip the access key prefix
+ assertArrayEquals(
+ "foobar".getBytes(StandardCharsets.US_ASCII),
+ AwsAccountIdentity.decodeBase32("AKIAMZXW6YTBOI", 4, 14));
+ }
+
+ @ParameterizedTest(name = "base32 rejects {0}")
+ @ValueSource(strings = {"MZXW6YTBO1", "MZXW6YTBO0", "MZXW6YTBO8", "MZXW6YTBO=", "MZXW6YTBO-"})
+ void base32RejectsCharactersOutsideTheAlphabet(String value) {
+ assertNull(AwsAccountIdentity.decodeBase32(value, 0, value.length()));
+ }
+
+ @Test
+ void extractsAccountFromEncodedBytes() {
+ // marker | (123456789012 << 7) | 0x55 = 0x8E5F_4C8D_0A55 big-endian over 6 bytes; the low 7
+ // bits and the trailing bytes are not part of the account and are ignored.
+ byte[] bytes = {(byte) 0x8E, 0x5F, 0x4C, (byte) 0x8D, 0x0A, 0x55, 0x12, 0x34, 0x56, 0x78};
+ assertEquals("123456789012", AwsAccountIdentity.accountFromEncodedBytes(bytes));
+ byte[] one = {(byte) 0x80, 0, 0, 0, 0, (byte) 0x80};
+ assertEquals("000000000001", AwsAccountIdentity.accountFromEncodedBytes(one));
+ assertNull(AwsAccountIdentity.accountFromEncodedBytes(new byte[5]));
+ }
+
+ @Test
+ void rejectsBytesWithoutTheFormatMarker() {
+ // Same account bits as above but bit 47 clear: a legacy key, not an encoded account.
+ byte[] legacy = {0x0E, 0x5F, 0x4C, (byte) 0x8D, 0x0A, 0x55, 0x12, 0x34, 0x56, 0x78};
+ assertNull(AwsAccountIdentity.accountFromEncodedBytes(legacy));
+ }
+
+ @ParameterizedTest(name = "rejects legacy access key id {0}")
+ @ValueSource(
+ strings = {
+ // AWS documentation example key: fifth character I is below Q, so pre-2019 format
+ "AKIAIOSFODNN7EXAMPLE",
+ "AKIAJ7EXAMPLEEXAMPLE",
+ "ASIAAAAAAAAAAAAAAAAA",
+ "AKIAPZZZZZZZZZZZZZZZ",
+ })
+ void rejectsLegacyAccessKeyIds(String value) {
+ assertNull(AwsAccountIdentity.accountFromAccessKeyId(value));
+ }
+
+ private static String syntheticKey(String prefix, long account) {
+ // format marker (bit 47) set, low 7 bits are not part of the account
+ long firstSixBytes = AwsAccountIdentity.ACCOUNT_FORMAT_MARKER | (account << 7) | 0x55L;
+ byte[] raw = new byte[10];
+ for (int i = 5; i >= 0; i--) {
+ raw[i] = (byte) (firstSixBytes & 0xff);
+ firstSixBytes >>>= 8;
+ }
+ raw[6] = (byte) 0x12;
+ raw[7] = (byte) 0x34;
+ raw[8] = (byte) 0x56;
+ raw[9] = (byte) 0x78;
+ String alphabet = "ABCDEFGHIJKLMNOPQRSTUVWXYZ234567";
+ StringBuilder out = new StringBuilder(prefix);
+ int buffer = 0;
+ int bits = 0;
+ for (byte b : raw) {
+ buffer = (buffer << 8) | (b & 0xff);
+ bits += 8;
+ while (bits >= 5) {
+ out.append(alphabet.charAt((buffer >> (bits - 5)) & 0x1f));
+ bits -= 5;
+ }
+ }
+ assertEquals(20, out.length());
+ return out.toString();
+ }
+}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-common/src/test/java/datadog/trace/instrumentation/aws/AwsArnTest.java b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/test/java/datadog/trace/instrumentation/aws/AwsArnTest.java
new file mode 100644
index 00000000000..306f115f4bb
--- /dev/null
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-common/src/test/java/datadog/trace/instrumentation/aws/AwsArnTest.java
@@ -0,0 +1,91 @@
+package datadog.trace.instrumentation.aws;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.CsvSource;
+import org.junit.jupiter.params.provider.NullAndEmptySource;
+import org.junit.jupiter.params.provider.ValueSource;
+
+class AwsArnTest {
+
+ @ParameterizedTest(name = "parses {0}")
+ @CsvSource(
+ nullValues = "NULL",
+ value = {
+ "arn:aws:dynamodb:us-east-1:123456789012:table/orders, aws, dynamodb, us-east-1, 123456789012, table/orders",
+ "arn:aws:dynamodb:us-east-1:123456789012:table/orders/index/by-user, aws, dynamodb, us-east-1, 123456789012, table/orders/index/by-user",
+ "arn:aws-cn:dynamodb:cn-north-1:123456789012:table/orders, aws-cn, dynamodb, cn-north-1, 123456789012, table/orders",
+ "arn:aws-us-gov:sns:us-gov-west-1:123456789012:alerts, aws-us-gov, sns, us-gov-west-1, 123456789012, alerts",
+ "arn:aws:s3:::my-bucket, aws, s3, NULL, NULL, my-bucket",
+ "arn:aws:iam::123456789012:role/app, aws, iam, NULL, 123456789012, role/app",
+ "arn:aws:states:us-east-1:123456789012:stateMachine:orders:extra:colons, aws, states, us-east-1, 123456789012, stateMachine:orders:extra:colons",
+ })
+ void parsesFields(
+ String value,
+ String partition,
+ String service,
+ String region,
+ String account,
+ String resource) {
+ AwsArn arn = AwsArn.parse(value);
+ assertNotNull(arn);
+ assertEquals(value, arn.raw());
+ assertEquals(partition, arn.partition());
+ assertEquals(service, arn.service());
+ assertEquals(region, arn.region());
+ assertEquals(account, arn.account());
+ assertEquals(resource, arn.resource());
+ }
+
+ @ParameterizedTest(name = "rejects {0}")
+ @NullAndEmptySource
+ @ValueSource(
+ strings = {
+ "orders",
+ "arn:",
+ "arn:aws",
+ "arn:aws:dynamodb",
+ "arn:aws:dynamodb:us-east-1",
+ "arn:aws:dynamodb:us-east-1:123456789012",
+ "arn:aws:dynamodb:us-east-1:123456789012:",
+ "arn::dynamodb:us-east-1:123456789012:table/orders",
+ "arn:aws::us-east-1:123456789012:table/orders",
+ "table/orders",
+ "https://sqs.us-east-1.amazonaws.com/123456789012/q",
+ })
+ void rejectsNonArn(String value) {
+ assertNull(AwsArn.parse(value));
+ }
+
+ @ParameterizedTest(name = "account of {0} is null")
+ @ValueSource(
+ strings = {
+ "arn:aws:dynamodb:us-east-1:12345:table/orders",
+ "arn:aws:dynamodb:us-east-1:12345678901a:table/orders",
+ "arn:aws:dynamodb:us-east-1::table/orders",
+ })
+ void accountMustBeTwelveDigits(String value) {
+ AwsArn arn = AwsArn.parse(value);
+ assertNotNull(arn);
+ assertNull(arn.account());
+ }
+
+ @ParameterizedTest(name = "table name of {0} is {1}")
+ @CsvSource(
+ nullValues = "NULL",
+ value = {
+ "arn:aws:dynamodb:us-east-1:123456789012:table/orders, orders",
+ "arn:aws:dynamodb:us-east-1:123456789012:table/orders/index/by-user, orders",
+ "arn:aws:dynamodb:us-east-1:123456789012:table/orders/stream/2024-01-01T00:00:00.000, orders",
+ "arn:aws:dynamodb:us-east-1:123456789012:table/, NULL",
+ "arn:aws:dynamodb:us-east-1:123456789012:backup/orders, NULL",
+ "arn:aws:sns:us-east-1:123456789012:table/orders, orders",
+ "arn:aws:s3:::my-bucket, NULL",
+ })
+ void dynamoDbTableName(String value, String expected) {
+ assertEquals(expected, AwsArn.parse(value).dynamoDbTableName());
+ }
+}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/build.gradle b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/build.gradle
index e8173d6afe0..610823222b1 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/build.gradle
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/build.gradle
@@ -38,6 +38,7 @@ addTestSuiteExtendingForDir('latestDsmForkedTest', 'latestDsmTest', 'dsmTest')
dependencies {
compileOnly group: 'com.amazonaws', name: 'aws-java-sdk-core', version: '1.11.0'
+ implementation project(':dd-java-agent:instrumentation:aws-java:aws-java-common')
// Include httpclient instrumentation for testing because it is a dependency for aws-sdk.
testImplementation project(':dd-java-agent:instrumentation:apache-httpclient:apache-httpclient-4.0')
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkClientDecorator.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkClientDecorator.java
index 07027e0f4f6..a94a10afa83 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkClientDecorator.java
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkClientDecorator.java
@@ -25,6 +25,8 @@
import datadog.trace.bootstrap.instrumentation.api.Tags;
import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString;
import datadog.trace.bootstrap.instrumentation.decorator.HttpClientDecorator;
+import datadog.trace.instrumentation.aws.AwsAccountIdentity;
+import datadog.trace.instrumentation.aws.AwsArn;
import java.net.URI;
import java.util.List;
import java.util.Locale;
@@ -179,12 +181,31 @@ protected void doOnRequest(final AgentSpan span, final Request request) {
}
String tableName = access.getTableName(originalRequest);
if (null != tableName) {
+ // A TableName given as an ARN carries the owning account; a bare name is resolved by
+ // DynamoDB in the requestor's own account. SDK v1 does not expose the signing credentials
+ // to request handlers, so only the ARN form yields an account here. Other services with a
+ // TableName member (Timestream, Keyspaces, ...) only get the plain table name tags.
+ AwsArn arn =
+ awsSimplifiedServiceName != null && awsSimplifiedServiceName.startsWith("dynamodb")
+ ? AwsArn.parse(tableName)
+ : null;
+ if (arn != null) {
+ String account = arn.account();
+ if (account != null) {
+ span.setTag(InstrumentationTags.AWS_ACCOUNT, account);
+ }
+ String bareName = arn.dynamoDbTableName();
+ if (bareName != null) {
+ // Only a table ARN is a table ARN; any other resource keeps the raw value as the name.
+ span.setTag(InstrumentationTags.AWS_TABLE_ARN, arn.raw());
+ tableName = bareName;
+ }
+ }
span.setTag(InstrumentationTags.AWS_TABLE_NAME, tableName);
span.setTag(InstrumentationTags.TABLE_NAME, tableName);
bestPrecursor = InstrumentationTags.AWS_TABLE_NAME;
bestPeerService = tableName;
}
-
// Set peer.service based on Config for serverless functions
if (Config.get().isAwsServerless()) {
URI uri = request.getEndpoint();
@@ -247,6 +268,23 @@ protected void doOnRequest(final AgentSpan span, final Request request) {
}
}
+ /**
+ * Tags derived from the request that are only trustworthy once the service accepted it. S3
+ * enforces {@code ExpectedBucketOwner} (403 on mismatch), so the asserted owner names the bucket
+ * owner only on a successful response.
+ */
+ public void onSuccessfulRequest(final AgentSpan span, final Request> request) {
+ if (!"s3".equalsIgnoreCase(simplifyServiceName(request.getServiceName()))) {
+ return;
+ }
+ final AmazonWebServiceRequest originalRequest = request.getOriginalRequest();
+ final String expectedBucketOwner =
+ GetterAccess.of(originalRequest).getExpectedBucketOwner(originalRequest);
+ if (AwsAccountIdentity.isAccountId(expectedBucketOwner)) {
+ span.setTag(InstrumentationTags.AWS_ACCOUNT, expectedBucketOwner);
+ }
+ }
+
public void onServiceResponse(
final AgentSpan span, final String awsService, final Response response) {
if ("s3".equalsIgnoreCase(simplifyServiceName(awsService))
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkModule.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkModule.java
index 8c370ce03b8..29db032dfc0 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkModule.java
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/AwsSdkModule.java
@@ -30,6 +30,8 @@ public String[] helperClassNames() {
packageName + ".TracingRequestHandler",
packageName + ".AwsNameCache",
packageName + ".OnErrorDecorator",
+ "datadog.trace.instrumentation.aws.AwsAccountIdentity",
+ "datadog.trace.instrumentation.aws.AwsArn",
};
}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/GetterAccess.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/GetterAccess.java
index ed87408c1bd..37e69fa18ca 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/GetterAccess.java
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/GetterAccess.java
@@ -39,6 +39,7 @@ static GetterAccess of(final Object request) {
private final MethodHandle getPublishBatchRequestEntries;
private final MethodHandle getApproximateArrivalTimestamp;
private final MethodHandle getTableName;
+ private final MethodHandle getExpectedBucketOwner;
private GetterAccess(final Class> objectType) {
operationName =
@@ -55,6 +56,7 @@ private GetterAccess(final Class> objectType) {
getApproximateArrivalTimestamp =
findGetter(objectType, "getApproximateArrivalTimestamp", Date.class);
getTableName = findStringGetter(objectType, "getTableName");
+ getExpectedBucketOwner = findStringGetter(objectType, "getExpectedBucketOwner");
}
String getOperationNameFromType() {
@@ -101,6 +103,10 @@ String getTableName(final Object object) {
return invokeForString(getTableName, object);
}
+ String getExpectedBucketOwner(final Object object) {
+ return invokeForString(getExpectedBucketOwner, object);
+ }
+
Date getApproximateArrivalTimestamp(final Object object) {
return invoke(getApproximateArrivalTimestamp, object);
}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/TracingRequestHandler.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/TracingRequestHandler.java
index 84e6f7ad142..afa65f9cd32 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/TracingRequestHandler.java
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/main/java/datadog/trace/instrumentation/aws/v0/TracingRequestHandler.java
@@ -96,6 +96,7 @@ public void afterResponse(final Request> request, final Response> response)
span = AgentSpan.fromContext(context);
if (span != null) {
DECORATE.onResponse(span, response);
+ DECORATE.onSuccessfulRequest(span, request);
DECORATE.onServiceResponse(span, request.getServiceName(), response);
DECORATE.beforeFinish(span);
span.finish();
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/test/groovy/AWS1ClientTest.groovy b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/test/groovy/AWS1ClientTest.groovy
index c5ea6ca23cf..c0b241748d2 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/test/groovy/AWS1ClientTest.groovy
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-1.11/src/test/groovy/AWS1ClientTest.groovy
@@ -210,6 +210,7 @@ abstract class AWS1ClientTest extends VersionedNamingTestBase {
"S3" | "GetObject" | "GET" | "/someBucket/someKey" | AmazonS3ClientBuilder.standard().withPathStyleAccessEnabled(true).withEndpointConfiguration(endpoint).withCredentials(credentialsProvider).build() | { c -> c.getObject("someBucket", "someKey") } | ["aws.bucket.name": "someBucket", "bucketname": "someBucket", "aws.object.key": "someKey"] | "" | "aws.bucket.name" | null
"DynamoDBv2" | "CreateTable" | "POST" | "/" | AmazonDynamoDBClientBuilder.standard().withEndpointConfiguration(endpoint).withCredentials(credentialsProvider).build() | { c -> c.createTable(new CreateTableRequest("sometable", null)) } | ["aws.table.name": "sometable", "tablename": "sometable"] | "" | "aws.table.name" | null
"DynamoDBv2" | "GetItem" | "POST" | "/" | AmazonDynamoDBClientBuilder.standard().withEndpointConfiguration(endpoint).withCredentials(credentialsProvider).build() | { c -> c.getItem(new GetItemRequest("sometable", ["attribute": new AttributeValue("somevalue")])) } | ["aws.table.name": "sometable", "tablename": "sometable"] | "" | "aws.table.name" | null
+ "DynamoDBv2" | "GetItem" | "POST" | "/" | AmazonDynamoDBClientBuilder.standard().withEndpointConfiguration(endpoint).withCredentials(credentialsProvider).build() | { c -> c.getItem(new GetItemRequest("arn:aws:dynamodb:us-east-1:123456789012:table/sometable", ["attribute": new AttributeValue("somevalue")])) } | ["aws.table.name": "sometable", "tablename": "sometable", "aws.table.arn": "arn:aws:dynamodb:us-east-1:123456789012:table/sometable", "aws_account": "123456789012"] | "" | "aws.table.name" | null
"Kinesis" | "DeleteStream" | "POST" | "/" | AmazonKinesisClientBuilder.standard().withEndpointConfiguration(endpoint).withCredentials(credentialsProvider).build() | { c -> c.deleteStream(new DeleteStreamRequest().withStreamName("somestream")) } | ["aws.stream.name": "somestream", "streamname": "somestream"] | "" | "aws.stream.name" | null
"SQS" | "CreateQueue" | "POST" | "/" | AmazonSQSClientBuilder.standard().withEndpointConfiguration(endpoint).withCredentials(credentialsProvider).build() | { c -> c.createQueue(new CreateQueueRequest("somequeue")) } | ["aws.queue.name": "somequeue", "queuename": "somequeue"] | """
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/build.gradle b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/build.gradle
index cfdc266d8ec..bd1a87ab6c6 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/build.gradle
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/build.gradle
@@ -36,7 +36,7 @@ def fixedSdkVersion = '2.20.33' // 2.20.34 is missing and breaks IDEA import
dependencies {
compileOnly group: 'software.amazon.awssdk', name: 'aws-core', version: '2.2.0'
payloadTaggingTestContainerImage(image('localstack/localstack:4.2.0', 'test.localstack.image'))
- testImplementation project(':dd-java-agent:instrumentation:aws-java:aws-java-common')
+ implementation project(':dd-java-agent:instrumentation:aws-java:aws-java-common')
// Include httpclient instrumentation for testing because it is a dependency for aws-sdk.
testImplementation project(':dd-java-agent:instrumentation:apache-httpclient:apache-httpclient-4.0')
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkClientDecorator.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkClientDecorator.java
index 7c1a1a85b90..f6082dc4127 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkClientDecorator.java
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkClientDecorator.java
@@ -28,7 +28,12 @@
import datadog.trace.bootstrap.instrumentation.api.Tags;
import datadog.trace.bootstrap.instrumentation.api.UTF8BytesString;
import datadog.trace.bootstrap.instrumentation.decorator.HttpClientDecorator;
+import datadog.trace.instrumentation.aws.AwsAccountIdentity;
+import datadog.trace.instrumentation.aws.AwsArn;
import datadog.trace.payloadtags.PayloadTagsData;
+import java.lang.invoke.MethodHandle;
+import java.lang.invoke.MethodHandles;
+import java.lang.invoke.MethodType;
import java.net.URI;
import java.time.Instant;
import java.util.ArrayList;
@@ -41,6 +46,9 @@
import java.util.Set;
import javax.annotation.Nonnull;
import javax.annotation.ParametersAreNonnullByDefault;
+import software.amazon.awssdk.auth.credentials.AwsCredentials;
+import software.amazon.awssdk.auth.signer.AwsSignerExecutionAttribute;
+import software.amazon.awssdk.awscore.AwsExecutionAttribute;
import software.amazon.awssdk.awscore.AwsResponse;
import software.amazon.awssdk.core.SdkBytes;
import software.amazon.awssdk.core.SdkField;
@@ -52,6 +60,7 @@
import software.amazon.awssdk.core.interceptor.SdkExecutionAttribute;
import software.amazon.awssdk.http.SdkHttpRequest;
import software.amazon.awssdk.http.SdkHttpResponse;
+import software.amazon.awssdk.regions.Region;
public class AwsSdkClientDecorator extends HttpClientDecorator
implements CarrierSetter {
@@ -96,6 +105,14 @@ public class AwsSdkClientDecorator extends HttpClientDecorator new ExecutionAttribute<>("KinesisStreamArn"));
+ // ExpectedBucketOwner is carried from the request to the response so the account is only tagged
+ // once S3 has accepted it (a mismatch is rejected with 403).
+ public static final ExecutionAttribute EXPECTED_BUCKET_OWNER_ATTRIBUTE =
+ InstanceStore.of(ExecutionAttribute.class)
+ .getOrCreate(
+ "DatadogExpectedBucketOwner",
+ () -> new ExecutionAttribute<>("DatadogExpectedBucketOwner"));
+
// not static because this object would be ClassLoader specific if multiple SDK instances were
// loaded by different loaders
private SdkField kinesisApproximateArrivalTimestampField = null;
@@ -135,6 +152,12 @@ public void onSdkRequest(
// S3
request.getValueForField("Bucket", String.class).ifPresent(name -> setBucketName(span, name));
if ("s3".equalsIgnoreCase(awsServiceName)) {
+ // S3 enforces ExpectedBucketOwner (403 on mismatch), so it names the owner only once the
+ // request has succeeded. Remember it here and tag it from the successful response.
+ request
+ .getValueForField("ExpectedBucketOwner", String.class)
+ .filter(AwsAccountIdentity::isAccountId)
+ .ifPresent(owner -> attributes.putAttribute(EXPECTED_BUCKET_OWNER_ATTRIBUTE, owner));
// gate "Key" extraction to S3 — DynamoDB's Key is Map, would CCE
request.getValueForField("Key", String.class).ifPresent(key -> setObjectKey(span, key));
if (traceConfig().isDataStreamsEnabled()) {
@@ -185,8 +208,18 @@ public void onSdkRequest(
}
});
- // DynamoDB
- request.getValueForField("TableName", String.class).ifPresent(name -> setTableName(span, name));
+ // DynamoDB. Other services (Timestream, Keyspaces, Athena, ...) also have TableName members;
+ // only DynamoDB resolves a bare name in the caller's account, so only it gets the ownership
+ // enrichment. Everything else keeps the plain table name tags.
+ if ("dynamodb".equalsIgnoreCase(awsServiceName)) {
+ request
+ .getValueForField("TableName", String.class)
+ .ifPresent(name -> onDynamoDbTable(span, name, attributes));
+ } else {
+ request
+ .getValueForField("TableName", String.class)
+ .ifPresent(name -> setTableName(span, name));
+ }
// DSM
if (traceConfig().isDataStreamsEnabled()) {
@@ -274,6 +307,105 @@ private static void setPeerService(
}
}
+ /**
+ * Tags the table plus its owning account and ARN. A TableName given as an ARN carries both. A
+ * bare name is, by DynamoDB's documented contract, resolved in the requestor's own account, so
+ * the account owning the signing credentials is the table owner.
+ */
+ private static void onDynamoDbTable(
+ final AgentSpan span, final String tableName, final ExecutionAttributes attributes) {
+ String name = tableName;
+ String account;
+ String tableArn = null;
+ AwsArn arn = AwsArn.parse(tableName);
+ if (arn != null) {
+ account = arn.account();
+ String bareName = arn.dynamoDbTableName();
+ if (bareName != null) {
+ // Only a table ARN is a table ARN; any other resource keeps the raw value as the name.
+ name = bareName;
+ tableArn = arn.raw();
+ }
+ } else {
+ account = callerAccount(attributes);
+ tableArn = AwsAccountIdentity.dynamoDbTableArn(regionOf(attributes), account, name);
+ }
+ setTableName(span, name);
+ if (account != null) {
+ span.setTag(InstrumentationTags.AWS_ACCOUNT, account);
+ }
+ if (tableArn != null) {
+ span.setTag(InstrumentationTags.AWS_TABLE_ARN, tableArn);
+ }
+ }
+
+ private static String regionOf(final ExecutionAttributes attributes) {
+ Region region = attributes.getAttribute(AwsExecutionAttribute.AWS_REGION);
+ return region == null ? null : region.id();
+ }
+
+ /**
+ * Account owning the credentials that sign this request. Read from the credentials object when
+ * the SDK exposes it (AwsCredentialsIdentity.accountId(), SDK 2.26+, populated by the STS, SSO,
+ * profile, process and container providers). Older SDKs can opt in to decoding it from the access
+ * key ID.
+ */
+ private static String callerAccount(final ExecutionAttributes attributes) {
+ AwsCredentials credentials =
+ attributes.getAttribute(AwsSignerExecutionAttribute.AWS_CREDENTIALS);
+ if (credentials == null) {
+ return null;
+ }
+ String account = credentialsAccountId(credentials);
+ if (account == null && Config.get().isAwsAccountFromAccessKeyEnabled()) {
+ account = AwsAccountIdentity.accountFromAccessKeyId(credentials.accessKeyId());
+ }
+ return account;
+ }
+
+ // Optional accountId() was introduced on AwsCredentialsIdentity after 2.2.0, so it is
+ // looked up reflectively per credentials class. A missing method is cached as this sentinel.
+ private static final MethodHandle NO_ACCOUNT_ID_GETTER =
+ MethodHandles.constant(Optional.class, Optional.empty());
+ private static final DDCache, MethodHandle> ACCOUNT_ID_GETTERS =
+ DDCaches.newFixedSizeCache(8);
+
+ private static MethodHandle accountIdGetter(final Class> type) {
+ try {
+ return MethodHandles.publicLookup()
+ .findVirtual(type, "accountId", MethodType.methodType(Optional.class));
+ } catch (NoSuchMethodException | IllegalAccessException e) {
+ return NO_ACCOUNT_ID_GETTER;
+ }
+ }
+
+ private static String credentialsAccountId(final AwsCredentials credentials) {
+ MethodHandle getter =
+ ACCOUNT_ID_GETTERS.computeIfAbsent(
+ credentials.getClass(), AwsSdkClientDecorator::accountIdGetter);
+ if (getter == NO_ACCOUNT_ID_GETTER) {
+ return null;
+ }
+ try {
+ Object value = getter.invoke(credentials);
+ if (value instanceof Optional) {
+ Object account = ((Optional>) value).orElse(null);
+ if (account instanceof String && AwsAccountIdentity.isAccountId((String) account)) {
+ return (String) account;
+ }
+ }
+ } catch (Throwable ignored) {
+ // treat as absent
+ }
+ return null;
+ }
+
+ private static void setBucketOwner(AgentSpan span, String owner) {
+ // aws_account only: the Agent's credit card obfuscator redacts 12-digit values under keys it
+ // does not know, and aws_account is on its allow list.
+ span.setTag(InstrumentationTags.AWS_ACCOUNT, owner);
+ }
+
private static void setBucketName(AgentSpan span, String name) {
span.setTag(InstrumentationTags.AWS_BUCKET_NAME, name);
span.setTag(InstrumentationTags.BUCKET_NAME, name);
@@ -317,6 +449,13 @@ public void onSdkResponse(
final AgentSpan span = fromContext(context);
Config config = Config.get();
String serviceName = attributes.getAttribute(SdkExecutionAttribute.SERVICE_NAME);
+
+ // Only a successful response proves the asserted ExpectedBucketOwner is the bucket owner.
+ String expectedBucketOwner = attributes.getAttribute(EXPECTED_BUCKET_OWNER_ATTRIBUTE);
+ if (expectedBucketOwner != null && httpResponse != null && httpResponse.isSuccessful()) {
+ setBucketOwner(span, expectedBucketOwner);
+ }
+
if (config.isCloudResponsePayloadTaggingEnabled()
&& config.isCloudPayloadTaggingEnabledFor(serviceName)) {
awsPojoToTags(span, ConfigDefaults.DEFAULT_TRACE_CLOUD_PAYLOAD_RESPONSE_TAG, response);
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkModule.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkModule.java
index bfe441c52e9..1a4f26fec58 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkModule.java
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/main/java/datadog/trace/instrumentation/aws/v2/AwsSdkModule.java
@@ -20,7 +20,9 @@ public AwsSdkModule() {
public String[] helperClassNames() {
return new String[] {
"datadog.trace.instrumentation.aws.v2.AwsSdkClientDecorator",
- "datadog.trace.instrumentation.aws.v2.TracingExecutionInterceptor"
+ "datadog.trace.instrumentation.aws.v2.TracingExecutionInterceptor",
+ "datadog.trace.instrumentation.aws.AwsAccountIdentity",
+ "datadog.trace.instrumentation.aws.AwsArn"
};
}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/payloadTaggingTest/java/datadog/trace/instrumentation/aws/v2/S3BucketOwnerForkedTest.java b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/payloadTaggingTest/java/datadog/trace/instrumentation/aws/v2/S3BucketOwnerForkedTest.java
new file mode 100644
index 00000000000..0b7af15b693
--- /dev/null
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/payloadTaggingTest/java/datadog/trace/instrumentation/aws/v2/S3BucketOwnerForkedTest.java
@@ -0,0 +1,117 @@
+package datadog.trace.instrumentation.aws.v2;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+
+import com.sun.net.httpserver.HttpServer;
+import datadog.trace.agent.test.AbstractInstrumentationTest;
+import datadog.trace.core.DDSpan;
+import java.io.IOException;
+import java.net.InetSocketAddress;
+import java.net.URI;
+import java.util.List;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import software.amazon.awssdk.auth.credentials.AwsBasicCredentials;
+import software.amazon.awssdk.auth.credentials.StaticCredentialsProvider;
+import software.amazon.awssdk.regions.Region;
+import software.amazon.awssdk.services.s3.S3Client;
+import software.amazon.awssdk.services.s3.model.GetObjectRequest;
+import software.amazon.awssdk.services.s3.model.S3Exception;
+
+/**
+ * Lives in this source set because its S3 model (2.18.40) has {@code ExpectedBucketOwner}; the base
+ * test suite pins s3 2.2.0 which predates the field.
+ */
+class S3BucketOwnerForkedTest extends AbstractInstrumentationTest {
+
+ private static HttpServer server;
+ private static S3Client client;
+
+ @BeforeAll
+ static void startServer() throws IOException {
+ server = HttpServer.create(new InetSocketAddress("localhost", 0), 0);
+ server.createContext(
+ "/",
+ exchange -> {
+ // Keys under "forbidden" simulate S3 rejecting the asserted ExpectedBucketOwner.
+ int status = exchange.getRequestURI().getPath().contains("forbidden") ? 403 : 200;
+ exchange.sendResponseHeaders(status, -1);
+ exchange.close();
+ });
+ server.start();
+ client =
+ S3Client.builder()
+ .endpointOverride(URI.create("http://localhost:" + server.getAddress().getPort()))
+ .region(Region.US_EAST_1)
+ .credentialsProvider(
+ StaticCredentialsProvider.create(AwsBasicCredentials.create("test", "test")))
+ .build();
+ }
+
+ @AfterAll
+ static void stopServer() {
+ client.close();
+ server.stop(0);
+ }
+
+ @Test
+ void expectedBucketOwnerTagsTheOwningAccount() throws Exception {
+ client.getObject(
+ GetObjectRequest.builder()
+ .bucket("somebucket")
+ .key("somekey")
+ .expectedBucketOwner("123456789012")
+ .build());
+
+ DDSpan span = firstSpan();
+ assertEquals("somebucket", span.getTag("aws.bucket.name"));
+ assertEquals("123456789012", span.getTag("aws_account"));
+ }
+
+ @Test
+ void noExpectedBucketOwnerMeansNoAccountTag() throws Exception {
+ client.getObject(GetObjectRequest.builder().bucket("somebucket").key("somekey").build());
+
+ DDSpan span = firstSpan();
+ assertEquals("somebucket", span.getTag("aws.bucket.name"));
+ assertFalse(span.getTags().containsKey("aws_account"));
+ }
+
+ @Test
+ void rejectedExpectedBucketOwnerIsNotTagged() throws Exception {
+ try {
+ client.getObject(
+ GetObjectRequest.builder()
+ .bucket("somebucket")
+ .key("forbidden")
+ .expectedBucketOwner("123456789012")
+ .build());
+ } catch (S3Exception expected) {
+ // S3 rejects a wrong owner with 403; the asserted owner must not be tagged as the account
+ }
+
+ DDSpan span = firstSpan();
+ assertEquals("somebucket", span.getTag("aws.bucket.name"));
+ assertFalse(span.getTags().containsKey("aws_account"));
+ }
+
+ @Test
+ void malformedExpectedBucketOwnerIsIgnored() throws Exception {
+ client.getObject(
+ GetObjectRequest.builder()
+ .bucket("somebucket")
+ .key("somekey")
+ .expectedBucketOwner("not-an-account")
+ .build());
+
+ assertFalse(firstSpan().getTags().containsKey("aws_account"));
+ }
+
+ private static DDSpan firstSpan() throws Exception {
+ writer.waitForTraces(1);
+ List trace = writer.firstTrace();
+ return trace.get(0);
+ }
+}
diff --git a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/test/groovy/Aws2ClientTest.groovy b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/test/groovy/Aws2ClientTest.groovy
index aa7ecf0cd24..5fd4871197a 100644
--- a/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/test/groovy/Aws2ClientTest.groovy
+++ b/dd-java-agent/instrumentation/aws-java/aws-java-sdk-2.2/src/test/groovy/Aws2ClientTest.groovy
@@ -433,6 +433,77 @@ abstract class Aws2ClientTest extends VersionedNamingTestBase {
"""
}
+ def "DynamoDb request with a table ARN tags the owning account and table ARN"() {
+ setup:
+ def client = DynamoDbClient.builder()
+ .endpointOverride(server.address)
+ .region(Region.AP_NORTHEAST_1)
+ .credentialsProvider(CREDENTIALS_PROVIDER)
+ .build()
+ responseBody.set("")
+ def tableArn = "arn:aws:dynamodb:ap-northeast-1:123456789012:table/sometable"
+
+ when:
+ client.getItem(GetItemRequest.builder().tableName(tableArn).key(["attribute": AttributeValue.builder().s("somevalue").build()]).build())
+ TEST_WRITER.waitForTraces(1)
+
+ then:
+ assertTraces(1) {
+ trace(1) {
+ span {
+ serviceName expectedService("DynamoDb", "GetItem")
+ operationName expectedOperation("DynamoDb", "GetItem")
+ resourceName "DynamoDb.GetItem"
+ spanType DDSpanTypes.HTTP_CLIENT
+ errored false
+ measured true
+ parent()
+ tags {
+ "$Tags.COMPONENT" "java-aws-sdk"
+ "$Tags.SPAN_KIND" Tags.SPAN_KIND_CLIENT
+ "$Tags.PEER_HOSTNAME" "localhost"
+ "$Tags.PEER_PORT" server.address.port
+ "$Tags.HTTP_METHOD" "POST"
+ "$Tags.HTTP_STATUS" 200
+ "aws.service" "DynamoDb"
+ "aws_service" "DynamoDb"
+ "aws.operation" "GetItem"
+ "aws.agent" "java-aws-sdk"
+ "aws.requestId" "UNKNOWN"
+ // the bare name is what peer.service and the *name tags carry
+ "aws.table.name" "sometable"
+ "tablename" "sometable"
+ "aws.table.arn" tableArn
+ "aws_account" "123456789012"
+ peerServiceFrom("aws.table.name")
+ urlTags("${server.address}/", ExpectedQueryParams.getExpectedQueryParams("GetItem"))
+ defaultTags(false, true)
+ }
+ }
+ }
+ }
+ }
+
+ def "DynamoDb request with a bare table name has no account when the credentials carry none"() {
+ setup:
+ def client = DynamoDbClient.builder()
+ .endpointOverride(server.address)
+ .region(Region.AP_NORTHEAST_1)
+ .credentialsProvider(CREDENTIALS_PROVIDER)
+ .build()
+ responseBody.set("")
+
+ when:
+ client.getItem(GetItemRequest.builder().tableName("sometable").key(["attribute": AttributeValue.builder().s("somevalue").build()]).build())
+ TEST_WRITER.waitForTraces(1)
+
+ then:
+ def tags = TEST_WRITER.firstTrace().first().tags
+ tags["aws.table.name"] == "sometable"
+ !tags.containsKey("aws_account")
+ !tags.containsKey("aws.table.arn")
+ }
+
def "timeout and retry errors captured"() {
setup:
def server = httpServer {
diff --git a/dd-trace-api/src/main/java/datadog/trace/api/config/TraceInstrumentationConfig.java b/dd-trace-api/src/main/java/datadog/trace/api/config/TraceInstrumentationConfig.java
index 137e9805519..d94c7e63108 100644
--- a/dd-trace-api/src/main/java/datadog/trace/api/config/TraceInstrumentationConfig.java
+++ b/dd-trace-api/src/main/java/datadog/trace/api/config/TraceInstrumentationConfig.java
@@ -205,6 +205,8 @@ public final class TraceInstrumentationConfig {
public static final String AXIS_PROMOTE_RESOURCE_NAME = "trace.axis.promote.resource-name";
public static final String SQS_BODY_PROPAGATION_ENABLED = "trace.sqs.body.propagation.enabled";
+ public static final String AWS_ACCOUNT_FROM_ACCESS_KEY_ENABLED =
+ "trace.aws.account.from.access.key.enabled";
public static final String TRACE_RESOURCE_RENAMING_ENABLED = "trace.resource.renaming.enabled";
diff --git a/internal-api/src/main/java/datadog/trace/api/Config.java b/internal-api/src/main/java/datadog/trace/api/Config.java
index 2db203dcbd4..1c71a49cac9 100644
--- a/internal-api/src/main/java/datadog/trace/api/Config.java
+++ b/internal-api/src/main/java/datadog/trace/api/Config.java
@@ -581,6 +581,7 @@
import static datadog.trace.api.config.RumConfig.RUM_TRACK_RESOURCES;
import static datadog.trace.api.config.RumConfig.RUM_TRACK_USER_INTERACTION;
import static datadog.trace.api.config.RumConfig.RUM_VERSION;
+import static datadog.trace.api.config.TraceInstrumentationConfig.AWS_ACCOUNT_FROM_ACCESS_KEY_ENABLED;
import static datadog.trace.api.config.TraceInstrumentationConfig.AXIS_PROMOTE_RESOURCE_NAME;
import static datadog.trace.api.config.TraceInstrumentationConfig.CASSANDRA_KEYSPACE_STATEMENT_EXTRACTION_ENABLED;
import static datadog.trace.api.config.TraceInstrumentationConfig.CODE_ORIGIN_FOR_SPANS_ENABLED;
@@ -1316,6 +1317,7 @@ public static String getHostName() {
private final boolean awsPropagationEnabled;
private final boolean sqsPropagationEnabled;
private final boolean sqsBodyPropagationEnabled;
+ private final boolean awsAccountFromAccessKeyEnabled;
private final boolean kafkaClientPropagationEnabled;
private final Set kafkaClientPropagationDisabledTopics;
@@ -3143,6 +3145,8 @@ PROFILING_DATADOG_PROFILER_ENABLED, isDatadogProfilerSafeInCurrentEnvironment())
awsPropagationEnabled = isPropagationEnabled(true, "aws", "aws-sdk");
sqsPropagationEnabled = isPropagationEnabled(true, "sqs");
sqsBodyPropagationEnabled = configProvider.getBoolean(SQS_BODY_PROPAGATION_ENABLED, false);
+ awsAccountFromAccessKeyEnabled =
+ configProvider.getBoolean(AWS_ACCOUNT_FROM_ACCESS_KEY_ENABLED, false);
kafkaClientPropagationEnabled = isPropagationEnabled(true, "kafka", "kafka.client");
kafkaClientPropagationDisabledTopics =
@@ -5127,6 +5131,10 @@ public boolean isSqsBodyPropagationEnabled() {
return sqsBodyPropagationEnabled;
}
+ public boolean isAwsAccountFromAccessKeyEnabled() {
+ return awsAccountFromAccessKeyEnabled;
+ }
+
public boolean isKafkaClientPropagationEnabled() {
return kafkaClientPropagationEnabled;
}
@@ -6946,6 +6954,8 @@ public String toString() {
+ debuggerCodeOriginEnabled
+ ", awsPropagationEnabled="
+ awsPropagationEnabled
+ + ", awsAccountFromAccessKeyEnabled="
+ + awsAccountFromAccessKeyEnabled
+ ", sqsPropagationEnabled="
+ sqsPropagationEnabled
+ ", kafkaClientPropagationEnabled="
diff --git a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/InstrumentationTags.java b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/InstrumentationTags.java
index 0c1054e7776..5a84543e55d 100644
--- a/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/InstrumentationTags.java
+++ b/internal-api/src/main/java/datadog/trace/bootstrap/instrumentation/api/InstrumentationTags.java
@@ -37,6 +37,11 @@ public class InstrumentationTags {
public static final String AWS_REQUEST_ID = "aws.requestId";
public static final String AWS_STORAGE_CLASS = "aws.storage.class";
+ // Owning account of the addressed resource. aws_account matches the tag dd-trace-py sets and the
+ // dimension tag on the AWS integration metrics.
+ public static final String AWS_ACCOUNT = "aws_account";
+ public static final String AWS_TABLE_ARN = "aws.table.arn";
+
// These are temporary keys used for span pointer hash calculation
public static final String S3_ETAG = "s3.eTag";
public static final String DYNAMO_PRIMARY_KEY_1 = "dynamodb.primary_key_1";
diff --git a/metadata/supported-configurations.json b/metadata/supported-configurations.json
index 6bc49c2679c..a04ea6b98b8 100644
--- a/metadata/supported-configurations.json
+++ b/metadata/supported-configurations.json
@@ -4940,6 +4940,14 @@
"aliases": ["DD_AWSADD_SPAN_POINTERS"]
}
],
+ "DD_TRACE_AWS_ACCOUNT_FROM_ACCESS_KEY_ENABLED": [
+ {
+ "version": "A",
+ "type": "boolean",
+ "default": "false",
+ "aliases": []
+ }
+ ],
"DD_TRACE_AWS_DYNAMODB_ENABLED": [
{
"version": "A",