diff --git a/CHANGELOG.md b/CHANGELOG.md index bbe0e8c..e52f26d 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 fb9d71a..ba81f4b 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 e90acee..b013b99 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 c4307e5..50da134 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 738eb1c..a4a1616 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 4a1c551..75659d8 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 3cd68ff..48e0107 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 1da5ecf..2395c5e 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 6af97a3..5e0d19a 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 b5bfb16..42fa7ee 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 38a6145..8244223 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 ffc2dd6..9539047 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 fe263a4..b970285 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 59b34c5..444fbb0 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 93bf7ba..fc60b08 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 4a35373..f7e8289 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 f18f68d..9e06bcf 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 5858acd..e724ee7 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 33fc868..bd29dc1 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 201f2f9..91d9383 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 0ca214d..086189d 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.