From 2c8a5c718fbedd929c5afefb9d1e89d977b104a9 Mon Sep 17 00:00:00 2001 From: Scott Friedman <3011922+scttfrdmn@users.noreply.github.com> Date: Thu, 24 Sep 2026 18:02:52 -0400 Subject: [PATCH] fix(#856): the shared UUID-shaped identifier, and the messaging/storage family MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Substrate had two shared draw sites, not one. randomHex was the known one; the other was a UUID-shaped generator declared in the Lambda plugin that six other services publish an identifier from — an ECS task id, a Step Functions execution name, an SQS message id, an EventBridge event id, a CloudWatch Logs upload sequence token and a Service Quotas request id. Tiering per service family would have left it on crypto/rand until the last of the six moved, so it moves here with Lambda's own revision ids. Its rendering is unchanged: sixteen derived bytes in UUID *shape*, 8-4-4-4-12 lowercase hex without the RFC 4122 version and variant bits. IDMint.UUID sets those bits, and using it here would have changed which bytes a caller sees, which is not what #856 is about. Moving with it: SQS receipt handles, SNS subscription ids and notification-envelope message ids, EFS file-system/access-point/mount-target ids, FSx file-system ids and Lustre mount names, and Transfer Family server ids. A send mints a message's initial receipt handle and each receive replaces it, matching real SQS, so the handle derives from the mint's ordinal rather than from the message — deriving from the message would make every receive of one message agree. Three helpers minted without a request context in scope and now take the mint as a parameter: buildSNSEnvelope, requestServiceQuotaIncrease and FSx's Lustre mount-name branch. 29 draw sites remain on crypto/rand, down from 32; the TODO on randomHex and on IDMint carries the new count. Tests: a SendMessageBatch whose three entries carry identical bodies mints three distinct message ids, which is the ordinal's test through the shared generator; and a recorded stream that sends and receives an SQS message, subscribes to an SNS topic and creates an EFS file system with an access point — each later request naming what an earlier one minted — replays with zero differences and StateValid true. Reverting generateSQSReceiptHandle alone to crypto/rand fails that test on the ReceiptHandle difference and on every state hash after it, so it is not vacuous. Two branches the change touched and nothing covered gained tests of their own: an unnamed StartExecution, and a Lustre deployment type other than SCRATCH_2. Refs #856. --- CHANGELOG.md | 23 ++- docs/services.md | 19 ++- docs/testing-guide.md | 7 +- emulator/cloudwatchlogs_plugin.go | 4 +- emulator/ec2_types.go | 2 +- emulator/ecs_plugin.go | 2 +- emulator/efs_plugin.go | 24 +-- emulator/eventbridge_plugin.go | 2 +- emulator/fsx_plugin.go | 12 +- emulator/fsx_plugin_test.go | 38 +++++ emulator/ids.go | 2 +- emulator/ids_test.go | 215 ++++++++++++++++++++++++++ emulator/lambda_plugin.go | 30 ++-- emulator/servicequotas_plugin.go | 6 +- emulator/sns_plugin.go | 17 +- emulator/sns_types.go | 6 +- emulator/sqs_plugin.go | 41 ++--- emulator/stepfunctions_plugin.go | 4 +- emulator/stepfunctions_plugin_test.go | 45 ++++++ emulator/transfer_plugin.go | 2 +- emulator/transfer_types.go | 10 +- 21 files changed, 427 insertions(+), 84 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index bbe0e8c1..e52f26d9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -169,9 +169,30 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 `StateValid` true, twice over and in a second process. `generateRequestID`, which is the seed, gained a random suffix: it was the wall clock alone, and every other identifier now derives from it, so two requests served within one tick of a coarse clock would have minted the same volume ID. - Remaining on `crypto/rand`: 32 draw sites in the other services, migrating one family at a time on + Remaining on `crypto/rand`: 29 draw sites in the other services, migrating one family at a time on #856; EC2 key-pair material, which needs a deterministic reader into the key generator rather than a string; and substrate's own event, snapshot and replay IDs, which no AWS call observes. +- **The shared UUID-shaped identifier, and the messaging/storage family, are derived too** (#856). + Substrate had two shared draw sites, not one. `randomHex` was the known one; the other was a + UUID-shaped generator declared in the Lambda plugin that **six other services** publish an + identifier from — an ECS task ID, a Step Functions execution name, an SQS message ID, an + EventBridge event ID, a CloudWatch Logs upload sequence token and a Service Quotas request ID — so + it moved to `IDMint` together with Lambda's own revision IDs rather than waiting for each of those + services' turn in the per-family tiering. Its rendering is unchanged: sixteen derived bytes in + UUID *shape*, `8-4-4-4-12` lowercase hex without the RFC 4122 version and variant bits, because + #856 is about reproducing an identifier across a replay and not about changing which bytes a + caller sees. Moving with it: SQS receipt handles, SNS subscription IDs and notification-envelope + message IDs, EFS file-system/access-point/mount-target IDs, FSx file-system IDs and Lustre mount + names, and Transfer Family server IDs. A send mints a message's initial receipt handle and each + receive replaces it, which matches real SQS — two `ReceiveMessage` calls returning one message + hand back two handles, and only the most recent deletes — so the handle derives from the mint's + ordinal rather than from the message, and each call's handle replays as the one that call + recorded. Three helpers that minted + without a request context in scope now take the mint as a parameter: `buildSNSEnvelope`, + `requestServiceQuotaIncrease` and FSx's Lustre mount-name branch. A recorded stream that sends and + receives an SQS message, subscribes to an SNS topic and creates an EFS file system with an access + point — each later request naming what an earlier one minted — now replays with **zero** + differences and `StateValid` true. - **A stream recorded under a seed replays under the same seed** (#1140). Every seedable outcome in substrate is written through a control-plane endpoint, and only the AWS path recorded anything — so a seed never entered the event stream. A replay opens by resetting the whole `StateManager`, and a diff --git a/docs/services.md b/docs/services.md index fb9d71ab..ba81f4b2 100644 --- a/docs/services.md +++ b/docs/services.md @@ -2206,9 +2206,22 @@ Three kinds of value stay random, and one more is still migrating: ID, a NAT gateway's private IP from its gateway ID, a secret's ARN from its name, and CloudFormation's [stack and change-set ARNs](#stack-and-change-set-arns-are-deterministic), which predate this rule and are what generalising it was modelled on. -- EC2, IAM and STS identifiers are derived today. The remaining services are migrating one family - at a time, tracked on #856; until a service moves, its identifiers are still drawn from - `crypto/rand` and a replay of a stream creating one of its resources still diverges. +- EC2, IAM, STS, SQS, SNS, Lambda, EFS, FSx, Transfer, ECS, Step Functions, EventBridge, + CloudWatch Logs and Service Quotas identifiers are derived today. The remaining services are + migrating one family at a time, tracked on #856; until a service moves, its identifiers are + still drawn from `crypto/rand` and a replay of a stream creating one of its resources still + diverges. + +Six of those services publish an identifier from one shared generator rather than declaring their +own, so they moved together: an ECS task ID, a Step Functions execution name, an SQS message ID, +an EventBridge event ID, a CloudWatch Logs upload sequence token and a Service Quotas request ID +are all the same sixteen derived bytes rendered in UUID *shape* — `8-4-4-4-12` lowercase hex +without the RFC 4122 version and variant bits, which is the form substrate published before it +derived them and is unchanged by deriving them. + +An SQS send mints a message's initial receipt handle and each receive replaces it, matching real +SQS: two `ReceiveMessage` calls that return the same message hand back different handles and only +the most recent one deletes. A replay of either call reproduces the handle that call recorded. --- diff --git a/docs/testing-guide.md b/docs/testing-guide.md index e90aceeb..b013b99f 100644 --- a/docs/testing-guide.md +++ b/docs/testing-guide.md @@ -335,9 +335,10 @@ makes `Differences` empty — and `StateValid` true, with `WithRecordedStateHash all; before #856 every such stream diverged on its first create, which is why the #1140 tests above are built on caller-chosen bucket and key names instead. -Two caveats. **Only EC2, IAM and STS identifiers are derived so far**; the remaining -services are migrating one family at a time, and until a service moves, a replay of a -stream creating one of its resources still diverges. And a recording made against an +Two caveats. **Not every service's identifiers are derived yet.** EC2, IAM, STS, SQS, SNS, +Lambda, EFS, FSx, Transfer, ECS, Step Functions, EventBridge, CloudWatch Logs and Service +Quotas are; the rest are migrating one family at a time, and until a service moves, a +replay of a stream creating one of its resources still diverges. And a recording made against an **unfrozen** clock can still diverge on a `state_hash_after` even when every identifier matches, because a handler reading the live clock stamps its record a few hundred nanoseconds after the event's own timestamp — invisible in a response rendering seconds, diff --git a/emulator/cloudwatchlogs_plugin.go b/emulator/cloudwatchlogs_plugin.go index c4307e52..50da1349 100644 --- a/emulator/cloudwatchlogs_plugin.go +++ b/emulator/cloudwatchlogs_plugin.go @@ -408,7 +408,7 @@ func (p *CloudWatchLogsPlugin) createLogStream(ctx *RequestContext, req *AWSRequ LogStreamName: body.LogStreamName, ARN: cwLogStreamARN(ctx.Region, ctx.AccountID, body.LogGroupName, body.LogStreamName), CreationTime: now, - UploadSequenceToken: generateLambdaRevisionID(), + UploadSequenceToken: generateLambdaRevisionID(ctx.IDs), } data, err := json.Marshal(ls) if err != nil { @@ -596,7 +596,7 @@ func (p *CloudWatchLogsPlugin) putLogEvents(ctx *RequestContext, req *AWSRequest var ls CWLogStream if json.Unmarshal(streamData, &ls) == nil { ls.LastIngestionTime = now - ls.UploadSequenceToken = generateLambdaRevisionID() + ls.UploadSequenceToken = generateLambdaRevisionID(ctx.IDs) if updated, marshalErr := json.Marshal(ls); marshalErr == nil { _ = p.state.Put(goCtx, cloudwatchLogsNamespace, streamKey, updated) } diff --git a/emulator/ec2_types.go b/emulator/ec2_types.go index 738eb1c9..a4a16166 100644 --- a/emulator/ec2_types.go +++ b/emulator/ec2_types.go @@ -495,7 +495,7 @@ func generateAssociationID(m *IDMint) string { // flag day: a caller moves by taking a mint and calling [IDMint.Hex] with the same width. // EC2's own ids no longer come through here. // -// TODO(#856): 32 draw sites remain on crypto/rand, tiered by service family on the issue; +// TODO(#856): 29 draw sites remain on crypto/rand, tiered by service family on the issue; // delete this function when the last caller moves. func randomHex(n int) string { b := make([]byte, n) diff --git a/emulator/ecs_plugin.go b/emulator/ecs_plugin.go index 4a1c551e..75659d8b 100644 --- a/emulator/ecs_plugin.go +++ b/emulator/ecs_plugin.go @@ -830,7 +830,7 @@ func (p *ECSPlugin) runTask(ctx *RequestContext, req *AWSRequest) (*AWSResponse, var tasks []ECSTask for i := 0; i < body.Count; i++ { - taskID := generateLambdaRevisionID()[:16] + taskID := generateLambdaRevisionID(ctx.IDs)[:16] task := ECSTask{ TaskArn: fmt.Sprintf("arn:aws:ecs:%s:%s:task/%s/%s", ctx.Region, ctx.AccountID, clusterName, taskID), TaskDefinitionArn: td.TaskDefinitionArn, diff --git a/emulator/efs_plugin.go b/emulator/efs_plugin.go index 3cd68ffe..48e01076 100644 --- a/emulator/efs_plugin.go +++ b/emulator/efs_plugin.go @@ -94,7 +94,7 @@ func (p *EFSPlugin) createFileSystem(reqCtx *RequestContext, req *AWSRequest) (* input.ThroughputMode = "bursting" } - fsID := generateEFSFileSystemID() + fsID := generateEFSFileSystemID(reqCtx.IDs) arn := fmt.Sprintf("arn:aws:elasticfilesystem:%s:%s:file-system/%s", reqCtx.Region, reqCtx.AccountID, fsID) @@ -256,7 +256,7 @@ func (p *EFSPlugin) createAccessPoint(reqCtx *RequestContext, req *AWSRequest) ( return nil, efsBadRequest("FileSystemId is required") } - apID := generateEFSAccessPointID() + apID := generateEFSAccessPointID(reqCtx.IDs) arn := fmt.Sprintf("arn:aws:elasticfilesystem:%s:%s:access-point/%s", reqCtx.Region, reqCtx.AccountID, apID) @@ -380,7 +380,7 @@ func (p *EFSPlugin) createMountTarget(reqCtx *RequestContext, req *AWSRequest) ( return nil, efsBadRequest("FileSystemId is required") } - mtID := generateEFSMountTargetID() + mtID := generateEFSMountTargetID(reqCtx.IDs) mt := EFSMountTarget{ MountTargetID: mtID, FileSystemID: input.FileSystemID, @@ -695,17 +695,17 @@ func efsJSONResponse(status int, v interface{}) (*AWSResponse, error) { }, nil } -// generateEFSFileSystemID generates an EFS file system ID (fs- + 8 hex chars). -func generateEFSFileSystemID() string { - return "fs-" + randomHex(8) +// generateEFSFileSystemID mints an EFS file system ID (fs- + 16 hex chars) from m. +func generateEFSFileSystemID(m *IDMint) string { + return "fs-" + m.Hex(8) } -// generateEFSAccessPointID generates an EFS access point ID (fsap- + 8 hex chars). -func generateEFSAccessPointID() string { - return "fsap-" + randomHex(8) +// generateEFSAccessPointID mints an EFS access point ID (fsap- + 16 hex chars) from m. +func generateEFSAccessPointID(m *IDMint) string { + return "fsap-" + m.Hex(8) } -// generateEFSMountTargetID generates an EFS mount target ID (fsmt- + 8 hex chars). -func generateEFSMountTargetID() string { - return "fsmt-" + randomHex(8) +// generateEFSMountTargetID mints an EFS mount target ID (fsmt- + 16 hex chars) from m. +func generateEFSMountTargetID(m *IDMint) string { + return "fsmt-" + m.Hex(8) } diff --git a/emulator/eventbridge_plugin.go b/emulator/eventbridge_plugin.go index 1da5ecfa..2395c5e9 100644 --- a/emulator/eventbridge_plugin.go +++ b/emulator/eventbridge_plugin.go @@ -465,7 +465,7 @@ func (p *EventBridgePlugin) putEvents(ctx *RequestContext, req *AWSRequest) (*AW Detail: entry.Detail, EventBusName: busName, Time: now, - EventID: generateLambdaRevisionID(), + EventID: generateLambdaRevisionID(ctx.IDs), } existing = append(existing, ev) results = append(results, resultEntry{EventID: ev.EventID}) diff --git a/emulator/fsx_plugin.go b/emulator/fsx_plugin.go index 6af97a3e..5e0d19a4 100644 --- a/emulator/fsx_plugin.go +++ b/emulator/fsx_plugin.go @@ -56,10 +56,10 @@ type FSxTag struct { Value string `json:"Value"` } -// generateFSxFileSystemID returns a unique FSx file system ID of the form -// "fs-" followed by 8 lowercase hex digits. -func generateFSxFileSystemID() string { - return "fs-" + randomHex(8) +// generateFSxFileSystemID mints a unique FSx file system ID from m, of the form +// "fs-" followed by 16 lowercase hex digits. +func generateFSxFileSystemID(m *IDMint) string { + return "fs-" + m.Hex(8) } // fsxDNSName derives a DNS name for a file system based on its ID and region. @@ -141,7 +141,7 @@ func (p *FSxPlugin) createFileSystem(ctx *RequestContext, req *AWSRequest) (*AWS input.StorageType = "SSD" } - fsID := generateFSxFileSystemID() + fsID := generateFSxFileSystemID(ctx.IDs) arn := fmt.Sprintf("arn:aws:fsx:%s:%s:file-system/%s", ctx.Region, ctx.AccountID, fsID) // Derive VPC ID from the first subnet when available (simplified). @@ -171,7 +171,7 @@ func (p *FSxPlugin) createFileSystem(ctx *RequestContext, req *AWSRequest) (*AWS if lustreDeploymentType == "SCRATCH_2" || lustreDeploymentType == "" { lustreMountName = "fsx" } else { - lustreMountName = randomHex(8) + lustreMountName = ctx.IDs.Hex(8) } } diff --git a/emulator/fsx_plugin_test.go b/emulator/fsx_plugin_test.go index b5bfb167..42fa7ee2 100644 --- a/emulator/fsx_plugin_test.go +++ b/emulator/fsx_plugin_test.go @@ -339,3 +339,41 @@ func TestFSx_SDKTargetRouting(t *testing.T) { require.NoError(t, json.Unmarshal(body, &out)) assert.True(t, strings.HasPrefix(out.FileSystem.FileSystemId, "fs-")) } + +// TestFSx_LustreMountNameIsMintedForANonScratchDeployment covers the other half of the +// MountName branch: SCRATCH_2 always reports "fsx", and every other Lustre deployment type +// reports a minted value instead. +// +// Since #856 that value derives from the request's own ID rather than from crypto/rand, so a +// replayed CreateFileSystem reports the mount name its recording reported. The assertion is on +// the shape and on its difference from the SCRATCH_2 constant; the derivation itself is +// asserted end-to-end in ids_test.go. +func TestFSx_LustreMountNameIsMintedForANonScratchDeployment(t *testing.T) { + ts := httptest.NewServer(newFSxTestServer(t)) + t.Cleanup(ts.Close) + + mountNameFor := func(deploymentType string) string { + resp := fsxRequest(t, ts, "CreateFileSystem", `{ + "FileSystemType": "LUSTRE", + "StorageCapacity": 1200, + "SubnetIds": ["subnet-12345678"], + "LustreConfiguration": {"DeploymentType": "`+deploymentType+`"} + }`) + body := readFSxBody(t, resp) + require.Equal(t, http.StatusOK, resp.StatusCode, "%s", body) + var created struct { + FileSystem struct { + LustreConfiguration struct { + MountName string `json:"MountName"` + } `json:"LustreConfiguration"` + } `json:"FileSystem"` + } + require.NoError(t, json.Unmarshal(body, &created)) + return created.FileSystem.LustreConfiguration.MountName + } + + assert.Equal(t, "fsx", mountNameFor("SCRATCH_2"), + "SCRATCH_2's mount name is the documented constant, not a minted value") + assert.Regexp(t, `^[0-9a-f]{16}$`, mountNameFor("PERSISTENT_1"), + "a persistent deployment reports a minted mount name") +} diff --git a/emulator/ids.go b/emulator/ids.go index 38a61458..82442232 100644 --- a/emulator/ids.go +++ b/emulator/ids.go @@ -61,7 +61,7 @@ import ( // already derived from its inputs: a public IP from its instance id, a secret's ARN from its // name, CloudFormation's stack UUIDs from account and region. // -// TODO(#856): 32 draw sites remain on crypto/rand, tiered by service family on the issue. +// TODO(#856): 29 draw sites remain on crypto/rand, tiered by service family on the issue. // IDMint mints the identifiers one request publishes, derived from that request's own id so // that replaying the request mints the same ones. diff --git a/emulator/ids_test.go b/emulator/ids_test.go index ffc2dd65..9539047f 100644 --- a/emulator/ids_test.go +++ b/emulator/ids_test.go @@ -2,6 +2,7 @@ package emulator_test import ( "encoding/base64" + "encoding/json" "encoding/xml" "io" "net/http" @@ -439,3 +440,217 @@ func idsEC2Call(t *testing.T, ts *emulator.TestServer, params map[string]string) require.Equal(t, http.StatusOK, resp.StatusCode, "%s: %s", params["Action"], body) return body } + +// Tier 2 of #856: the shared UUID-shaped draw site, and the messaging/storage family. +// +// [emulator.IDMint] arrived with EC2, IAM and STS moved onto it. The assertions below are +// about the second of substrate's two shared draw sites — the one Lambda declares and six +// other services publish an identifier from — and about the SQS/SNS/EFS/FSx/Transfer group +// that moved with it. They are the same two layers as above: a distinctness property +// asserted directly on one request that mints several identifiers, and a replayed stream +// asserted over the wire. + +// TestIDs_OneSendMessageBatchMintsDistinctMessageIDs is the ordinal's test for the shared +// generator, in the same shape as the RunInstances one above: several identifiers from a +// single request, so a mint that ignored its ordinal would answer one ID three times. +// +// SQS is the interesting caller because it mints twice per message on two different paths — +// a message ID on send and a receipt handle on receive — so the two must not collide either. +func TestIDs_OneSendMessageBatchMintsDistinctMessageIDs(t *testing.T) { + t.Parallel() + ts := emulator.StartTestServer(t) + ts.FreezeTime() + + queueURL := idsSQSQueue(t, ts, "ids-batch") + + params := map[string]string{"Action": "SendMessageBatch", "QueueUrl": queueURL} + for i := 1; i <= 3; i++ { + params["SendMessageBatchRequestEntry."+strconv.Itoa(i)+".Id"] = "e" + strconv.Itoa(i) + params["SendMessageBatchRequestEntry."+strconv.Itoa(i)+".MessageBody"] = "same body" + } + var batch struct { + Entries []struct { + MessageID string `xml:"MessageId"` + } `xml:"SendMessageBatchResult>SendMessageBatchResultEntry"` + } + require.NoError(t, xml.Unmarshal(idsSQSCall(t, ts, params), &batch)) + require.Len(t, batch.Entries, 3) + + minted := make([]string, 0, 3) + for _, e := range batch.Entries { + require.NotEmpty(t, e.MessageID) + minted = append(minted, e.MessageID) + } + assert.Len(t, idsUnique(minted), 3, + "one SendMessageBatch mints one message ID per entry, and the bodies are identical: %v", + minted) +} + +// TestReplay_TheSharedMintReplaysWithTheIdentifiersItMinted is the wire-level assertion for +// Tier 2, and the same claim [TestReplay_ACreateReplaysWithTheIdentifiersItMinted] makes for +// EC2 and IAM: a stream whose later requests *name* the identifiers its earlier ones minted +// replays with no differences and reaches the recorded state. +// +// The interlocking is what makes it more than a comparison of five responses. A re-minted +// queue URL, subscription ARN, file-system ID or receipt handle turns the request that names +// it into a refusal rather than into a differently-worded success. +func TestReplay_TheSharedMintReplaysWithTheIdentifiersItMinted(t *testing.T) { + t.Parallel() + ts := emulator.StartTestServer(t, + emulator.WithRecordedBodies(), emulator.WithRecordedStateHashes()) + require.True(t, ts.Store().RecordsStateHashes(), "precondition: state_hash_after is compared") + + idsRecordSharedMintCreates(t, ts) + + results, err := replayEngineFor(ts, emulator.ReplayConfig{ValidateState: true}). + Replay(t.Context(), replayStreamID) + require.NoError(t, err) + + assert.Positive(t, results.TotalEvents, "the stream has to contain the creates") + assert.Equal(t, results.TotalEvents, results.SuccessEvents, + "every recorded request is re-executed and answers") + assert.Empty(t, results.Differences, + "an identifier from the shared mint replays as the one recorded: %s", + replayDifferenceSummary(results)) + assert.True(t, results.StateValid, + "and the state it reaches is the recorded state: %v", results.StateErrors) +} + +// idsRecordSharedMintCreates records the Tier 2 stream: a create per family that mints +// through the shared generator or through one of the messaging/storage generators, each +// followed by a request naming what the create minted. +func idsRecordSharedMintCreates(t *testing.T, ts *emulator.TestServer) { + t.Helper() + + // Frozen for the reason [idsRecordInterlockedCreates] gives: a replay pins the clock to + // the recorded event's timestamp, so a handler that stamps a record off a live clock + // diverges in the state hash for a reason that has nothing to do with an identifier. + ts.FreezeTime() + + // SQS: a message ID from the shared generator on send, a receipt handle on receive, and + // a delete that names the handle. + queueURL := idsSQSQueue(t, ts, "ids-shared-mint") + idsSQSCall(t, ts, map[string]string{ + "Action": "SendMessage", "QueueUrl": queueURL, "MessageBody": "one", + }) + var received struct { + Handles []string `xml:"ReceiveMessageResult>Message>ReceiptHandle"` + } + require.NoError(t, xml.Unmarshal(idsSQSCall(t, ts, map[string]string{ + "Action": "ReceiveMessage", "QueueUrl": queueURL, "MaxNumberOfMessages": "1", + }), &received)) + require.Len(t, received.Handles, 1, "the message just sent has to be received") + idsSQSCall(t, ts, map[string]string{ + "Action": "DeleteMessage", "QueueUrl": queueURL, "ReceiptHandle": received.Handles[0], + }) + + // SNS: a subscription ID, then a get that names the ARN carrying it. + var topic struct { + ARN string `xml:"CreateTopicResult>TopicArn"` + } + require.NoError(t, xml.Unmarshal(idsSNSCall(t, ts, map[string]string{ + "Action": "CreateTopic", "Name": "ids-shared-mint", + }), &topic)) + require.NotEmpty(t, topic.ARN) + + var sub struct { + ARN string `xml:"SubscribeResult>SubscriptionArn"` + } + require.NoError(t, xml.Unmarshal(idsSNSCall(t, ts, map[string]string{ + "Action": "Subscribe", "TopicArn": topic.ARN, + "Protocol": "email", "Endpoint": "nobody@example.invalid", + }), &sub)) + require.NotEmpty(t, sub.ARN) + idsSNSCall(t, ts, map[string]string{ + "Action": "GetSubscriptionAttributes", "SubscriptionArn": sub.ARN, + }) + + // EFS: a file-system ID, then an access point and a mount target that name it. + var fs struct { + ID string `json:"FileSystemId"` + } + require.NoError(t, json.Unmarshal(idsEFSCall(t, ts, http.MethodPost, "/2015-02-01/file-systems", + map[string]any{"CreationToken": "ids-shared-mint"}), &fs)) + require.NotEmpty(t, fs.ID) + idsEFSCall(t, ts, http.MethodPost, "/2015-02-01/access-points", + map[string]any{"ClientToken": "ids-ap", "FileSystemId": fs.ID}) + idsEFSCall(t, ts, http.MethodGet, "/2015-02-01/file-systems?FileSystemId="+fs.ID, nil) +} + +// idsSQSQueue creates one SQS queue and returns its URL. +func idsSQSQueue(t *testing.T, ts *emulator.TestServer, name string) string { + t.Helper() + + var created struct { + URL string `xml:"CreateQueueResult>QueueUrl"` + } + body := idsSQSCall(t, ts, map[string]string{"Action": "CreateQueue", "QueueName": name}) + require.NoError(t, xml.Unmarshal(body, &created)) + require.NotEmpty(t, created.URL, "CreateQueue returned no URL: %s", body) + return created.URL +} + +// idsSQSCall issues an unsigned SQS query-protocol request and returns the response body. +func idsSQSCall(t *testing.T, ts *emulator.TestServer, params map[string]string) []byte { + t.Helper() + return idsQueryCall(t, ts, "sqs.us-east-1.amazonaws.com", "2012-11-05", params) +} + +// idsSNSCall issues an unsigned SNS query-protocol request and returns the response body. +func idsSNSCall(t *testing.T, ts *emulator.TestServer, params map[string]string) []byte { + t.Helper() + return idsQueryCall(t, ts, "sns.us-east-1.amazonaws.com", "2010-03-31", params) +} + +// idsQueryCall issues one unsigned query-protocol request against host and requires a 200. +func idsQueryCall(t *testing.T, ts *emulator.TestServer, host, version string, + params map[string]string, +) []byte { + t.Helper() + + form := url.Values{} + form.Set("Version", version) + for k, v := range params { + form.Set(k, v) + } + + req, err := http.NewRequestWithContext(t.Context(), http.MethodPost, ts.URL+"/", + strings.NewReader(form.Encode())) + require.NoError(t, err) + req.Host = host + req.Header.Set("Content-Type", "application/x-www-form-urlencoded") + + resp, err := http.DefaultClient.Do(req) + require.NoError(t, err) + body, err := io.ReadAll(resp.Body) + require.NoError(t, err) + require.NoError(t, resp.Body.Close()) + require.Equal(t, http.StatusOK, resp.StatusCode, "%s: %s", params["Action"], body) + return body +} + +// idsEFSCall issues one EFS REST/JSON request and requires a 2xx. +func idsEFSCall(t *testing.T, ts *emulator.TestServer, method, path string, body any) []byte { + t.Helper() + + var reader io.Reader + if body != nil { + raw, err := json.Marshal(body) + require.NoError(t, err) + reader = strings.NewReader(string(raw)) + } + req, err := http.NewRequestWithContext(t.Context(), method, ts.URL+path, reader) + require.NoError(t, err) + req.Host = "elasticfilesystem.us-east-1.amazonaws.com" + if body != nil { + req.Header.Set("Content-Type", "application/json") + } + + resp, err := http.DefaultClient.Do(req) + require.NoError(t, err) + out, err := io.ReadAll(resp.Body) + require.NoError(t, err) + require.NoError(t, resp.Body.Close()) + require.Less(t, resp.StatusCode, 300, "%s %s: %s", method, path, out) + return out +} diff --git a/emulator/lambda_plugin.go b/emulator/lambda_plugin.go index fe263a4d..b970285a 100644 --- a/emulator/lambda_plugin.go +++ b/emulator/lambda_plugin.go @@ -2,7 +2,6 @@ package emulator import ( "context" - "crypto/rand" "crypto/sha256" "encoding/base64" "encoding/json" @@ -280,7 +279,7 @@ func (p *LambdaPlugin) createFunction(ctx *RequestContext, req *AWSRequest) (*AW Environment: body.Environment.Variables, CodeSize: 0, CodeSha256: "", - RevisionID: generateLambdaRevisionID(), + RevisionID: generateLambdaRevisionID(ctx.IDs), State: "Active", PackageType: pkgType, Architectures: archs, @@ -395,7 +394,7 @@ func (p *LambdaPlugin) updateFunctionCode(ctx *RequestContext, req *AWSRequest, // fresh random value on every update — a caller comparing it across two updates // is asking whether the code changed, and a random value answers "always". A // RevisionID is not a digest of anything and stays random. - fn.RevisionID = generateLambdaRevisionID() + fn.RevisionID = generateLambdaRevisionID(ctx.IDs) fn.LastModified = p.tc.Now() switch { @@ -481,7 +480,7 @@ func (p *LambdaPlugin) updateFunctionConfiguration(ctx *RequestContext, req *AWS if body.Environment.Variables != nil { fn.Environment = body.Environment.Variables } - fn.RevisionID = generateLambdaRevisionID() + fn.RevisionID = generateLambdaRevisionID(ctx.IDs) fn.LastModified = p.tc.Now() resp, err := p.saveFunctionAndRespond(ctx.AccountID, ctx.Region, fn, http.StatusOK) @@ -1128,11 +1127,22 @@ func (p *LambdaPlugin) sizeS3Package(bucket, key, versionID string) (int64, stri return obj.Size, strings.Trim(obj.ETag, `"`) } -// generateLambdaRevisionID generates a unique revision/code ID. -func generateLambdaRevisionID() string { - b := make([]byte, 16) - _, _ = rand.Read(b) - return fmt.Sprintf("%x-%x-%x-%x-%x", b[0:4], b[4:6], b[6:8], b[8:10], b[10:16]) +// generateLambdaRevisionID returns a UUID-shaped revision/code ID minted from m. +// +// This is the second shared draw site after [randomHex], and the wider one of the two: +// six other services publish an identifier from here — ECS task IDs, Step Functions +// execution names, SQS message IDs, EventBridge event IDs, CloudWatch Logs sequence +// tokens and Service Quotas request IDs — so it moved to [IDMint] together with Lambda's +// own revision IDs rather than waiting for each of those services' turn (#856). +// +// The rendering is the raw hex of sixteen bytes in UUID *shape*, not an RFC 4122 +// version-4 UUID with its version and variant bits set. That is what this function +// published before; #856 is about making an identifier reproducible across a replay, not +// about changing which bytes a caller sees, so [IDMint.UUID] is deliberately not used +// here. +func generateLambdaRevisionID(m *IDMint) string { + h := m.Hex(16) + return h[0:8] + "-" + h[8:12] + "-" + h[12:16] + "-" + h[16:20] + "-" + h[20:32] } // parseInt parses a string into an int. @@ -1262,7 +1272,7 @@ func (p *LambdaPlugin) createEventSourceMapping(ctx *RequestContext, req *AWSReq functionARN = "arn:aws:lambda:" + ctx.Region + ":" + ctx.AccountID + ":function:" + input.FunctionName } - uuid := generateLambdaRevisionID() + uuid := generateLambdaRevisionID(ctx.IDs) esm := &ESMConfig{ UUID: uuid, FunctionARN: functionARN, diff --git a/emulator/servicequotas_plugin.go b/emulator/servicequotas_plugin.go index 59b34c5a..444fbb02 100644 --- a/emulator/servicequotas_plugin.go +++ b/emulator/servicequotas_plugin.go @@ -56,7 +56,7 @@ func (p *ServiceQuotasPlugin) HandleRequest(reqCtx *RequestContext, req *AWSRequ case "GetAWSDefaultServiceQuota": return p.getAWSDefaultServiceQuota(req) case "RequestServiceQuotaIncrease": - return p.requestServiceQuotaIncrease(sqAccountID(reqCtx), req) + return p.requestServiceQuotaIncrease(reqCtx.IDs, sqAccountID(reqCtx), req) // The second name is not one AWS has; it is the name this handler shipped // under, kept as an alias so a fixture that drives the plugin by a hand-built // X-Amz-Target keeps working. See the handler's doc comment (#636). @@ -196,7 +196,7 @@ func (p *ServiceQuotasPlugin) getAWSDefaultServiceQuota(req *AWSRequest) (*AWSRe return p.getServiceQuota(req) } -func (p *ServiceQuotasPlugin) requestServiceQuotaIncrease(accountID string, req *AWSRequest) (*AWSResponse, error) { +func (p *ServiceQuotasPlugin) requestServiceQuotaIncrease(mint *IDMint, accountID string, req *AWSRequest) (*AWSResponse, error) { var input struct { ServiceCode string `json:"ServiceCode"` QuotaCode string `json:"QuotaCode"` @@ -209,7 +209,7 @@ func (p *ServiceQuotasPlugin) requestServiceQuotaIncrease(accountID string, req return nil, sqIllegalArgument("ServiceCode and QuotaCode are required") } - id := generateSQSMessageID() // reuse UUID generator + id := generateSQSMessageID(mint) // reuse UUID generator qi := &QuotaIncrease{ ID: id, diff --git a/emulator/sns_plugin.go b/emulator/sns_plugin.go index 93bf7ba9..fc60b083 100644 --- a/emulator/sns_plugin.go +++ b/emulator/sns_plugin.go @@ -567,7 +567,7 @@ func (p *SNSPlugin) subscribe(ctx *RequestContext, req *AWSRequest) (*AWSRespons return nil, err } topicName := target.Name - subID := generateSNSSubID() + subID := generateSNSSubID(ctx.IDs) // The subscription ARN is minted under the account and Region of the topic, because a // subscription ARN is the topic's ARN with the subscription's identifier appended — its account // segment is the topic's, not the subscriber's. The record and the two indexes below stay keyed by @@ -910,7 +910,7 @@ func (p *SNSPlugin) publish(ctx *RequestContext, req *AWSRequest) (*AWSResponse, return nil, err } - msgID := generateSNSSubID() + msgID := generateSNSSubID(ctx.IDs) for _, subARN := range subIDs { sub, loadErr := p.loadSub(goCtx, ctx.AccountID, ctx.Region, subARN) @@ -947,7 +947,7 @@ func (p *SNSPlugin) dispatchToSubscriber(ctx *RequestContext, sub *SNSSubscripti } switch sub.Protocol { case "sqs": - envelope := p.buildSNSEnvelope(sub, message, subject) + envelope := p.buildSNSEnvelope(ctx.IDs, sub, message, subject) envelopeBytes, _ := json.Marshal(envelope) queueURL := sub.Endpoint _, err := p.registry.RouteRequest(ctx, &AWSRequest{ @@ -995,7 +995,7 @@ func (p *SNSPlugin) dispatchToSubscriber(ctx *RequestContext, sub *SNSSubscripti p.logger.Warn("sns dispatch to lambda failed", "function", fnName, "err", err) } case "http", "https": - envelope := p.buildSNSEnvelope(sub, message, subject) + envelope := p.buildSNSEnvelope(ctx.IDs, sub, message, subject) envelopeBytes, _ := json.Marshal(envelope) httpReq, reqErr := http.NewRequest(http.MethodPost, sub.Endpoint, bytes.NewReader(envelopeBytes)) if reqErr != nil { @@ -1017,11 +1017,12 @@ func (p *SNSPlugin) dispatchToSubscriber(ctx *RequestContext, sub *SNSSubscripti } } -// buildSNSEnvelope wraps a message in the standard SNS notification JSON envelope. -func (p *SNSPlugin) buildSNSEnvelope(sub *SNSSubscription, message, subject string) map[string]interface{} { +// buildSNSEnvelope wraps a message in the standard SNS notification JSON envelope, +// minting the envelope's MessageId from m. +func (p *SNSPlugin) buildSNSEnvelope(m *IDMint, sub *SNSSubscription, message, subject string) map[string]interface{} { return map[string]interface{}{ "Type": "Notification", - "MessageId": generateSNSSubID(), + "MessageId": generateSNSSubID(m), "TopicArn": sub.TopicARN, "Subject": subject, "Message": message, @@ -1112,7 +1113,7 @@ func (p *SNSPlugin) publishBatch(ctx *RequestContext, req *AWSRequest) (*AWSResp } message := req.Params[fmt.Sprintf("PublishBatchRequestEntries.member.%d.Message", i)] batchSubject := req.Params[fmt.Sprintf("PublishBatchRequestEntries.member.%d.Subject", i)] - msgID := generateSNSSubID() + msgID := generateSNSSubID(ctx.IDs) for _, subARN := range subIDs { sub, loadErr := p.loadSub(goCtx, ctx.AccountID, ctx.Region, subARN) if loadErr != nil || sub == nil { diff --git a/emulator/sns_types.go b/emulator/sns_types.go index 4a353733..f7e82891 100644 --- a/emulator/sns_types.go +++ b/emulator/sns_types.go @@ -86,7 +86,7 @@ func snsSubscriptionARN(region, accountID, topicName, subID string) string { return fmt.Sprintf("arn:aws:sns:%s:%s:%s:%s", region, accountID, topicName, subID) } -// generateSNSSubID returns a short random ID for SNS subscriptions. -func generateSNSSubID() string { - return randomHex(8) +// generateSNSSubID mints a short ID for SNS subscriptions from m. +func generateSNSSubID(m *IDMint) string { + return m.Hex(8) } diff --git a/emulator/sqs_plugin.go b/emulator/sqs_plugin.go index f18f68d8..9e06bcfa 100644 --- a/emulator/sqs_plugin.go +++ b/emulator/sqs_plugin.go @@ -2,8 +2,7 @@ package emulator import ( "context" - "crypto/md5" //nolint:gosec // SQS MD5OfBody is defined by the protocol; not used for security. - "crypto/rand" + "crypto/md5" //nolint:gosec // SQS MD5OfBody is defined by the protocol; not used for security. "crypto/sha256" //nolint:gosec // SHA-256 used for content-based deduplication; not for security. "encoding/json" "encoding/xml" @@ -957,7 +956,7 @@ func (p *SQSPlugin) sendMessage(ctx *RequestContext, req *AWSRequest) (*AWSRespo }) } // Record this deduplication ID. - msgID := generateSQSMessageID() + msgID := generateSQSMessageID(ctx.IDs) p.recordFIFODedup(context.Background(), urlKey, dedupID, msgID, p.tc.Now()) md5Body := computeMD5(msgBody) @@ -965,7 +964,7 @@ func (p *SQSPlugin) sendMessage(ctx *RequestContext, req *AWSRequest) (*AWSRespo now := p.tc.Now() msg := &SQSMessage{ MessageID: msgID, - ReceiptHandle: generateSQSReceiptHandle(), + ReceiptHandle: generateSQSReceiptHandle(ctx.IDs), Body: msgBody, MD5OfBody: md5Body, Attributes: map[string]string{ @@ -1019,14 +1018,14 @@ func (p *SQSPlugin) sendMessage(ctx *RequestContext, req *AWSRequest) (*AWSRespo }) } - msgID := generateSQSMessageID() + msgID := generateSQSMessageID(ctx.IDs) md5Body := computeMD5(msgBody) md5Attrs := sqsMD5OfMessageAttributes(msgAttrs) now := p.tc.Now() msg := &SQSMessage{ MessageID: msgID, - ReceiptHandle: generateSQSReceiptHandle(), + ReceiptHandle: generateSQSReceiptHandle(ctx.IDs), Body: msgBody, MD5OfBody: md5Body, Attributes: map[string]string{ @@ -1155,11 +1154,11 @@ func (p *SQSPlugin) sendMessageBatch(ctx *RequestContext, req *AWSRequest) (*AWS failures = append(failures, sqsBatchFailure(entry.ID, awsErr)) continue } - msgID := generateSQSMessageID() + msgID := generateSQSMessageID(ctx.IDs) md5Body := computeMD5(entry.MessageBody) msg := &SQSMessage{ MessageID: msgID, - ReceiptHandle: generateSQSReceiptHandle(), + ReceiptHandle: generateSQSReceiptHandle(ctx.IDs), Body: entry.MessageBody, MD5OfBody: md5Body, Attributes: map[string]string{ @@ -1240,12 +1239,12 @@ func (p *SQSPlugin) sendMessageBatch(ctx *RequestContext, req *AWSRequest) (*AWS continue } - msgID := generateSQSMessageID() + msgID := generateSQSMessageID(ctx.IDs) md5Body := computeMD5(body) msg := &SQSMessage{ MessageID: msgID, - ReceiptHandle: generateSQSReceiptHandle(), + ReceiptHandle: generateSQSReceiptHandle(ctx.IDs), Body: body, MD5OfBody: md5Body, Attributes: map[string]string{ @@ -1436,7 +1435,7 @@ func (p *SQSPlugin) receiveMessage(ctx *RequestContext, req *AWSRequest) (*AWSRe } // Update receipt handle and visibility timeout. - newHandle := generateSQSReceiptHandle() + newHandle := generateSQSReceiptHandle(ctx.IDs) msg.ReceiptHandle = newHandle msg.VisibleAfter = now.Add(time.Duration(visTimeout) * time.Second) msg.ReceiveCount++ @@ -1835,16 +1834,20 @@ func getAttrOrDefault(attrs map[string]string, key, fallback string) string { return fallback } -// generateSQSMessageID generates a unique SQS message ID. -func generateSQSMessageID() string { - return generateLambdaRevisionID() // Reuse UUID-style generator. +// generateSQSMessageID mints a unique SQS message ID from m. +func generateSQSMessageID(m *IDMint) string { + return generateLambdaRevisionID(m) // Reuse UUID-style generator. } -// generateSQSReceiptHandle generates a unique receipt handle. -func generateSQSReceiptHandle() string { - b := make([]byte, 32) - _, _ = rand.Read(b) //nolint:gosec // Receipt handle just needs to be unique, not cryptographically secure. - return fmt.Sprintf("%x", b) +// generateSQSReceiptHandle mints a unique receipt handle from m. +// +// A send mints the message's initial handle and each receive replaces it, so two +// ReceiveMessage calls returning the same message hand back different handles — which is +// also true of real SQS, where a handle belongs to a receive and only the most recent one +// deletes. That is why the mint's ordinal rather than the message ID is the right source +// here: deriving from the message would make every receive of one message agree. +func generateSQSReceiptHandle(m *IDMint) string { + return m.Hex(32) } // computeMD5 computes the hex MD5 of s. diff --git a/emulator/stepfunctions_plugin.go b/emulator/stepfunctions_plugin.go index 5858acdd..e724ee7b 100644 --- a/emulator/stepfunctions_plugin.go +++ b/emulator/stepfunctions_plugin.go @@ -665,7 +665,7 @@ func (p *StepFunctionsPlugin) startExecution(ctx *RequestContext, req *AWSReques execName := input.Name if execName == "" { - execName = "exec-" + generateLambdaRevisionID()[:8] + execName = "exec-" + generateLambdaRevisionID(ctx.IDs)[:8] } // The execution belongs to the state machine, so its ARN, its record and its index entry are all @@ -760,7 +760,7 @@ func (p *StepFunctionsPlugin) startSyncExecution(ctx *RequestContext, req *AWSRe execName := input.Name if execName == "" { - execName = "sync-" + generateLambdaRevisionID()[:8] + execName = "sync-" + generateLambdaRevisionID(ctx.IDs)[:8] } execArn := fmt.Sprintf("arn:aws:states:%s:%s:express:%s:%s", target.Region, target.AccountID, target.Name, execName) diff --git a/emulator/stepfunctions_plugin_test.go b/emulator/stepfunctions_plugin_test.go index 33fc8687..bd29dc16 100644 --- a/emulator/stepfunctions_plugin_test.go +++ b/emulator/stepfunctions_plugin_test.go @@ -1035,3 +1035,48 @@ func TestSFN_Choice_IsNull(t *testing.T) { assert.Equal(t, "SUCCEEDED", d2["status"]) assert.Contains(t, d2["output"], "was-present") } + +// TestStepFunctions_StartExecutionMintsANameWhenNoneIsGiven covers the branch that names an +// execution for a caller who did not name one. `name` is optional on StartExecution, and the +// minted name is `exec-` followed by eight characters of the shared UUID-shaped identifier — +// which since #856 is derived from the request's own ID, so replaying the start reproduces the +// execution ARN the recording answered with rather than minting a new one. +// +// Two properties are asserted together because each is the other's control: two starts under +// one mint get different names, so the mint's ordinal is honored; and a second mint over the +// same seed reproduces the first name, so the source is the request ID and not a draw. +func TestStepFunctions_StartExecutionMintsANameWhenNoneIsGiven(t *testing.T) { + arnFor := func(mint *emulator.IDMint, machine string) []string { + p, ctx := setupStepFunctionsPlugin(t) + ctx.IDs = mint + _, err := p.HandleRequest(ctx, sfnRequest("CreateStateMachine", map[string]any{ + "name": machine, + "definition": testSMDefinition, + "roleArn": "arn:aws:iam::123456789012:role/sfn-role", + })) + require.NoError(t, err) + + arns := make([]string, 0, 2) + for range 2 { + resp, startErr := p.HandleRequest(ctx, sfnRequest("StartExecution", map[string]any{ + "stateMachineArn": "arn:aws:states:us-east-1:123456789012:stateMachine:" + machine, + })) + require.NoError(t, startErr) + require.Equal(t, 200, resp.StatusCode) + arn, ok := sfnBody(t, resp)["executionArn"].(string) + require.True(t, ok) + arns = append(arns, arn) + } + return arns + } + + first := arnFor(emulator.NewIDMint("req-mints-a-name"), "UnnamedExecSM") + assert.Regexp(t, `:execution:UnnamedExecSM:exec-[0-9a-f]{8}$`, first[0], + "an unnamed execution is named exec- plus eight hex characters") + assert.NotEqual(t, first[0], first[1], + "two unnamed starts in one request mint two names") + + second := arnFor(emulator.NewIDMint("req-mints-a-name"), "UnnamedExecSM") + assert.Equal(t, first, second, + "the name derives from the request ID, so the same ID mints the same name (#856)") +} diff --git a/emulator/transfer_plugin.go b/emulator/transfer_plugin.go index 201f2f95..91d93836 100644 --- a/emulator/transfer_plugin.go +++ b/emulator/transfer_plugin.go @@ -82,7 +82,7 @@ func (p *TransferPlugin) createServer(reqCtx *RequestContext, req *AWSRequest) ( input.EndpointType = "PUBLIC" } - serverID := generateTransferServerID() + serverID := generateTransferServerID(reqCtx.IDs) server := TransferServer{ ServerID: serverID, Arn: fmt.Sprintf("arn:aws:transfer:%s:%s:server/%s", reqCtx.Region, reqCtx.AccountID, serverID), diff --git a/emulator/transfer_types.go b/emulator/transfer_types.go index 0ca214db..086189da 100644 --- a/emulator/transfer_types.go +++ b/emulator/transfer_types.go @@ -1,8 +1,6 @@ package emulator import ( - "crypto/rand" - "encoding/hex" "time" ) @@ -63,12 +61,10 @@ type TransferTag struct { Value string `json:"Value"` } -// generateTransferServerID generates a server ID in the form s-{17 hex chars}, +// generateTransferServerID mints a server ID from m in the form s-{17 hex chars}, // matching the real AWS Transfer Family server ID format. -func generateTransferServerID() string { - b := make([]byte, 9) // 9 bytes → 18 hex chars; we use 17 - _, _ = rand.Read(b) - return "s-" + hex.EncodeToString(b)[:17] +func generateTransferServerID(m *IDMint) string { + return "s-" + m.Hex(9)[:17] // 9 bytes → 18 hex chars; we use 17. } // State key helpers.