Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 41 additions & 1 deletion providers/aws/lambda/aliases.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package lambda

import (
"context"
"maps"
"time"

cerrors "github.com/stackshy/cloudemu/v2/errors"
Expand All @@ -11,6 +12,9 @@ import (

// CreateAlias creates a new alias pointing to a specific function version.
func (m *Mock) CreateAlias(_ context.Context, cfg driver.AliasConfig) (*driver.Alias, error) {
m.mu.Lock()
defer m.mu.Unlock()

fd, ok := m.funcs.Get(cfg.FunctionName)
if !ok {
return nil, cerrors.Newf(cerrors.NotFound, "function %s not found", cfg.FunctionName)
Expand Down Expand Up @@ -53,6 +57,9 @@ func (m *Mock) CreateAlias(_ context.Context, cfg driver.AliasConfig) (*driver.A

// UpdateAlias updates an existing alias configuration.
func (m *Mock) UpdateAlias(_ context.Context, cfg driver.AliasConfig) (*driver.Alias, error) {
m.mu.Lock()
defer m.mu.Unlock()

fd, ok := m.funcs.Get(cfg.FunctionName)
if !ok {
return nil, cerrors.Newf(cerrors.NotFound, "function %s not found", cfg.FunctionName)
Expand Down Expand Up @@ -109,8 +116,15 @@ func (m *Mock) UpdateAlias(_ context.Context, cfg driver.AliasConfig) (*driver.A
return &result, nil
}

// DeleteAlias removes an alias from a function.
// DeleteAlias removes an alias from a function together with the state scoped
// to the alias qualifier: its resource-based policy, function URL config,
// provisioned concurrency config and event invoke config. Real Lambda keeps
// these as sub-resources of the alias ARN, so a later alias with the same name
// starts with none of them instead of inheriting the old grants.
func (m *Mock) DeleteAlias(_ context.Context, functionName, aliasName string) error {
m.mu.Lock()
defer m.mu.Unlock()

fd, ok := m.funcs.Get(functionName)
if !ok {
return cerrors.Newf(cerrors.NotFound, "function %s not found", functionName)
Expand All @@ -121,10 +135,36 @@ func (m *Mock) DeleteAlias(_ context.Context, functionName, aliasName string) er
}

fd.aliases.Delete(aliasName)
dropQualifierState(&fd, aliasName)
m.funcs.Set(functionName, fd)

return nil
}

// dropQualifierState removes everything keyed by a version or alias qualifier:
// the resource-based policy, function URL config, event invoke config and
// provisioned concurrency config. The maps are replaced rather than edited so a
// reader holding an earlier funcData copy is unaffected. Callers hold m.mu.
func dropQualifierState(fd *funcData, qualifier string) {
fd.policies = withoutKey(fd.policies, qualifier)
fd.urlConfigs = withoutKey(fd.urlConfigs, qualifier)
fd.eventInvokeConfigs = withoutKey(fd.eventInvokeConfigs, qualifier)
fd.provisionedConcurrencyConfigs = withoutKey(fd.provisionedConcurrencyConfigs, qualifier)
}

// withoutKey returns src unchanged when key is absent, otherwise a copy of src
// without key.
func withoutKey[V any](src map[string]V, key string) map[string]V {
if _, ok := src[key]; !ok {
return src
}

next := maps.Clone(src)
delete(next, key)

return next
}

// GetAlias retrieves a specific alias for a function.
func (m *Mock) GetAlias(_ context.Context, functionName, aliasName string) (*driver.Alias, error) {
fd, ok := m.funcs.Get(functionName)
Expand Down
242 changes: 242 additions & 0 deletions providers/aws/lambda/engine_lock_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,242 @@
package lambda

import (
"context"
"net/url"
"sync"
"sync/atomic"
"testing"
"time"

"github.com/stackshy/cloudemu/v2/config"
"github.com/stackshy/cloudemu/v2/errors"
"github.com/stackshy/cloudemu/v2/services/serverless/driver"
)

// unblockedRead is how long a read on another function may take while an
// engine call is in flight. The engine call itself blocks until released.
const unblockedRead = 100 * time.Millisecond

// gatedEngine is a FunctionEngine whose Deploy and Remove block until the test
// releases them, standing in for a slow real engine (a container build).
type gatedEngine struct {
started chan string
release chan struct{}
}

func newGatedEngine() *gatedEngine {
return &gatedEngine{started: make(chan string, 4), release: make(chan struct{})}
}

//nolint:gocritic // fn is the by-value DTO defined by the FunctionEngine contract
func (g *gatedEngine) Deploy(_ context.Context, fn config.FunctionDeployment) error {
g.started <- fn.Name
<-g.release

return nil
}

func (*gatedEngine) Invoke(context.Context, string, []byte) (config.FunctionResult, error) {
return config.FunctionResult{}, nil
}

func (g *gatedEngine) Remove(_ context.Context, name string) error {
g.started <- name
<-g.release

return nil
}

// otherFunction creates a code-less function (no engine call) with a policy
// and a function URL, and returns the URL host.
func otherFunction(t *testing.T, m *Mock) string {
t.Helper()

ctx := context.Background()

if _, err := m.CreateFunction(ctx, driver.FunctionConfig{Name: "other", Runtime: "python3.12", Handler: "h"}); err != nil {
t.Fatalf("CreateFunction(other): %v", err)
}

if err := m.AddPermission(ctx, "other", "", driver.PermissionStatement{
StatementID: "s", Action: "lambda:InvokeFunction", Principal: "*",
}); err != nil {
t.Fatalf("AddPermission(other): %v", err)
}

cfg, err := m.CreateFunctionURLConfig(ctx, driver.FunctionURLConfig{FunctionName: "other", AuthType: "NONE"})
if err != nil {
t.Fatalf("CreateFunctionURLConfig(other): %v", err)
}

u, err := url.Parse(cfg.FunctionURL)
if err != nil {
t.Fatalf("parse %s: %v", cfg.FunctionURL, err)
}

return u.Host
}

// assertReadsNotBlocked checks the reads the auth gate and URL dispatch make
// on every request return promptly on an unrelated function.
func assertReadsNotBlocked(t *testing.T, m *Mock, host string) {
t.Helper()

ctx := context.Background()
begin := time.Now()

if _, stmts, err := m.PolicyStatements(ctx, "other", ""); err != nil || len(stmts) != 1 {
t.Fatalf("PolicyStatements(other) = %v, %v", stmts, err)
}

if _, err := m.ResolveFunctionURL(ctx, host); err != nil {
t.Fatalf("ResolveFunctionURL: %v", err)
}

if _, err := m.GetPolicy(ctx, "other", ""); err != nil {
t.Fatalf("GetPolicy(other): %v", err)
}

if took := time.Since(begin); took > unblockedRead {
t.Fatalf("reads on another function took %v while an engine call was in flight", took)
}
}

func waitStarted(t *testing.T, g *gatedEngine, want string) {
t.Helper()

select {
case got := <-g.started:
if got != want {
t.Fatalf("engine call for %q, want %q", got, want)
}
case <-time.After(5 * time.Second):
t.Fatalf("engine call for %q never started", want)
}
}

func TestEngineDeployDoesNotHoldLock(t *testing.T) {
eng := newGatedEngine()
m := newEngineMock(eng)
host := otherFunction(t, m)
ctx := context.Background()

created := make(chan error, 1)

go func() {
_, err := m.CreateFunction(ctx, engineFuncConfig())
created <- err
}()

waitStarted(t, eng, "real-fn")
assertReadsNotBlocked(t, m, host)

if _, err := m.CreateFunction(ctx, engineFuncConfig()); !errors.IsAlreadyExists(err) {
t.Fatalf("CreateFunction of a name being created err = %v, want AlreadyExists", err)
}

eng.release <- struct{}{}

if err := <-created; err != nil {
t.Fatalf("CreateFunction: %v", err)
}

// Code update: the deploy runs unlocked, and other changes to the same
// function are refused until it finishes.
updated := make(chan error, 1)

go func() {
_, err := m.UpdateFunction(ctx, "real-fn", driver.FunctionConfig{Code: []byte("v2")})
updated <- err
}()

waitStarted(t, eng, "real-fn")
assertReadsNotBlocked(t, m, host)

if _, err := m.UpdateFunction(ctx, "real-fn", driver.FunctionConfig{Timeout: 9}); !errors.IsAlreadyExists(err) {
t.Fatalf("UpdateFunction during a code deploy err = %v, want AlreadyExists (ResourceConflict)", err)
}

if _, err := m.PublishVersion(ctx, "real-fn", ""); !errors.IsAlreadyExists(err) {
t.Fatalf("PublishVersion during a code deploy err = %v, want AlreadyExists (ResourceConflict)", err)
}

if err := m.TagFunction(ctx, "real-fn", map[string]string{"k": "v"}); !errors.IsAlreadyExists(err) {
t.Fatalf("TagFunction during a code deploy err = %v, want AlreadyExists (ResourceConflict)", err)
}

if err := m.AddPermission(ctx, "real-fn", "", driver.PermissionStatement{
StatementID: "during", Action: "lambda:InvokeFunction", Principal: "*",
}); err != nil {
t.Fatalf("AddPermission during a code deploy: %v", err)
}

eng.release <- struct{}{}

if err := <-updated; err != nil {
t.Fatalf("UpdateFunction(code): %v", err)
}

// The update was applied to the entry as it was after the deploy, so the
// grant added meanwhile survived.
if _, stmts, _ := m.PolicyStatements(ctx, "real-fn", ""); len(stmts) != 1 {
t.Fatalf("policy after the code update = %v, want the grant added during the deploy", stmts)
}

// Delete: the function is gone at once; the name stays reserved until the
// engine Remove returns.
deleted := make(chan error, 1)

go func() { deleted <- m.DeleteFunction(ctx, "real-fn") }()

waitStarted(t, eng, "real-fn")
assertReadsNotBlocked(t, m, host)

if _, err := m.GetFunction(ctx, "real-fn"); !errors.IsNotFound(err) {
t.Fatalf("GetFunction during the engine remove err = %v, want NotFound", err)
}

if _, err := m.CreateFunction(ctx, engineFuncConfig()); !errors.IsAlreadyExists(err) {
t.Fatalf("CreateFunction during the engine remove err = %v, want AlreadyExists", err)
}

eng.release <- struct{}{}

if err := <-deleted; err != nil {
t.Fatalf("DeleteFunction: %v", err)
}
}

func TestConcurrentCreateFunctionSameName(t *testing.T) {
const callers = 8

m := newTestMock()
start := make(chan struct{})

var (
wg sync.WaitGroup
ok atomic.Int32
)

for range callers {
wg.Add(1)

go func() {
defer wg.Done()
<-start

if _, err := m.CreateFunction(context.Background(), defaultFuncConfig()); err == nil {
ok.Add(1)
} else if !errors.IsAlreadyExists(err) {
t.Errorf("CreateFunction err = %v, want nil or AlreadyExists", err)
}
}()
}

close(start)
wg.Wait()

if got := ok.Load(); got != 1 {
t.Fatalf("%d CreateFunction calls succeeded, want exactly 1", got)
}
}
Loading
Loading