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
6 changes: 1 addition & 5 deletions cmd/ateapi/internal/actoridentity/actoridentity.go
Original file line number Diff line number Diff line change
Expand Up @@ -307,11 +307,7 @@ func authenticateAtelet(ctx context.Context) (*ateletCaller, error) {
// validateWorkerRef checks the reference to the Worker the certificate is
// minted for. Workers are global-scoped, so the reference carries no atespace.
func validateWorkerRef(worker *ateapipb.ObjectRef) error {
fldPath := field.NewPath("worker")
if worker == nil {
return field.Required(fldPath, "")
}
return resources.ValidateGlobalObjectRef(worker, fldPath).ToAggregate()
return resources.ValidateGlobalObjectRef(worker, field.NewPath("worker")).ToAggregate()
}

// authorizeActor resolves the actor from the authenticated worker and verifies
Expand Down
23 changes: 13 additions & 10 deletions cmd/ateapi/internal/actoridentity/actoridentity_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -212,14 +212,17 @@ func TestMintCertReadsThroughStaleWorkerCache(t *testing.T) {
if err != nil {
t.Fatalf("read seeded worker: %v", err)
}
if worker.Status == nil {
worker.Status = &ateapipb.WorkerStatus{}
}
worker.Status.Assignment = &ateapipb.ActorAssignment{
Actor: (resources.ActorRef{Atespace: testAtespace, Name: testActorName}).ToObjectRef(),
ActorUid: actor.GetMetadata().GetUid(),
}
if err := st.UpdateWorker(ctx, worker, worker.GetMetadata().GetVersion()); err != nil {
_, err = st.UpdateWorker(ctx, testWorkerName, store.PreconditionFrom(worker), func(toUpdate *ateapipb.Worker) error {
if toUpdate.Status == nil {
toUpdate.Status = &ateapipb.WorkerStatus{}
}
toUpdate.Status.Assignment = &ateapipb.ActorAssignment{
Actor: (resources.ActorRef{Atespace: testAtespace, Name: testActorName}).ToObjectRef(),
ActorUid: actor.GetMetadata().GetUid(),
}
return nil
})
if err != nil {
t.Fatalf("assign worker in store: %v", err)
}
}
Expand Down Expand Up @@ -273,7 +276,7 @@ func TestMintCertReadsThroughWorkerCacheMiss(t *testing.T) {
if workerInStore {
// Phase 2: register and assign the worker in the store only,
// after the cache stopped listening.
if err := st.CreateWorker(ctx, &ateapipb.Worker{
if _, err := st.CreateWorker(ctx, &ateapipb.Worker{
Metadata: &ateapipb.ResourceMetadata{Name: testWorkerName},
WorkerNamespace: testPodNS,
WorkerPool: testPool,
Expand Down Expand Up @@ -461,7 +464,7 @@ func seedActor(t *testing.T, ctx context.Context, st store.Interface, f actorFix
if f.unassigned {
worker.Status.Assignment = nil
}
if err := st.CreateWorker(ctx, worker); err != nil {
if _, err := st.CreateWorker(ctx, worker); err != nil {
t.Fatalf("seed worker: %v", err)
}
}
Expand Down
20 changes: 2 additions & 18 deletions cmd/ateapi/internal/controlapi/atespace.go
Original file line number Diff line number Diff line change
Expand Up @@ -92,15 +92,7 @@ func (s *RPCService) GetAtespace(ctx context.Context, req *ateapipb.GetAtespaceR

func validateGetAtespaceRequest(req *ateapipb.GetAtespaceRequest) field.ErrorList {
var fldPath *field.Path
var errs field.ErrorList

if val, fldPath := req.Atespace, fldPath.Child("atespace"); val == nil {
errs = append(errs, field.Required(fldPath, ""))
} else {
errs = append(errs, resources.ValidateGlobalObjectRef(val, fldPath)...)
}

return errs
return resources.ValidateGlobalObjectRef(req.GetAtespace(), fldPath.Child("atespace"))
}

func (s *RPCService) ListAtespaces(ctx context.Context, req *ateapipb.ListAtespacesRequest) (*ateapipb.ListAtespacesResponse, error) {
Expand Down Expand Up @@ -159,13 +151,5 @@ func (s *RPCService) DeleteAtespace(ctx context.Context, req *ateapipb.DeleteAte

func validateDeleteAtespaceRequest(req *ateapipb.DeleteAtespaceRequest) field.ErrorList {
var fldPath *field.Path
var errs field.ErrorList

if val, fldPath := req.Atespace, fldPath.Child("atespace"); val == nil {
errs = append(errs, field.Required(fldPath, ""))
} else {
errs = append(errs, resources.ValidateGlobalObjectRef(val, fldPath)...)
}

return errs
return resources.ValidateGlobalObjectRef(req.GetAtespace(), fldPath.Child("atespace"))
}
8 changes: 5 additions & 3 deletions cmd/ateapi/internal/controlapi/crash.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ type crashActorStore interface {
GetActor(ctx context.Context, actorRef resources.ActorRef) (*ateapipb.Actor, error)
UpdateActor(ctx context.Context, actorRef resources.ActorRef, precondition store.Precondition, mutate func(toUpdate *ateapipb.Actor) error) (*ateapipb.Actor, error)
GetWorker(ctx context.Context, name string) (*ateapipb.Worker, error)
UpdateWorker(ctx context.Context, worker *ateapipb.Worker, expectedVersion int64) error
UpdateWorker(ctx context.Context, name string, precondition store.Precondition, mutate func(toUpdate *ateapipb.Worker) error) (*ateapipb.Worker, error)
}

// releaseWorker clears the worker's assignment if it still points at the given
Expand Down Expand Up @@ -145,8 +145,10 @@ func releaseWorker(ctx context.Context, st crashActorStore, actor *ateapipb.Acto
return sandboxClass, nil
}

worker.Status.Assignment = nil
if err := st.UpdateWorker(ctx, worker, worker.GetMetadata().GetVersion()); err != nil {
if _, err := st.UpdateWorker(ctx, workerName, store.PreconditionFrom(worker), func(toUpdate *ateapipb.Worker) error {
toUpdate.Status.Assignment = nil
return nil
}); err != nil {
return sandboxClass, fmt.Errorf("while releasing worker: %w", err)
}
return sandboxClass, nil
Expand Down
10 changes: 5 additions & 5 deletions cmd/ateapi/internal/controlapi/crash_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,7 @@ func seedWorker(t *testing.T, ctx context.Context, st store.Interface, actorRef
}
}
}
if err := st.CreateWorker(ctx, worker); err != nil {
if _, err := st.CreateWorker(ctx, worker); err != nil {
t.Fatalf("seed worker: %v", err)
}
}
Expand Down Expand Up @@ -200,7 +200,7 @@ func TestCrashActor(t *testing.T) {
},
},
}
if err := st.CreateWorker(ctx, worker); err != nil {
if _, err := st.CreateWorker(ctx, worker); err != nil {
t.Fatalf("CreateWorker: %v", err)
}
},
Expand Down Expand Up @@ -432,7 +432,7 @@ func TestCrashActor_Metrics(t *testing.T) {
},
},
}
if err := st.CreateWorker(ctx, worker); err != nil {
if _, err := st.CreateWorker(ctx, worker); err != nil {
t.Fatalf("CreateWorker: %v", err)
}

Expand Down Expand Up @@ -513,8 +513,8 @@ type failingUpdateWorkerStore struct {
err error
}

func (f failingUpdateWorkerStore) UpdateWorker(context.Context, *ateapipb.Worker, int64) error {
return f.err
func (f failingUpdateWorkerStore) UpdateWorker(context.Context, string, store.Precondition, func(*ateapipb.Worker) error) (*ateapipb.Worker, error) {
return nil, f.err
}

// A transient failure releasing the worker must not move the actor to the
Expand Down
6 changes: 4 additions & 2 deletions cmd/ateapi/internal/controlapi/functionaltest/actor_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2454,8 +2454,10 @@ func TestResumeActor_CrashesIfAssignedWorkerIsDraining(t *testing.T) {
if err != nil {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's do in follow up PR, but can we add tests for the Worker methods to functionaltest/worker_test.go?

t.Fatalf("GetWorker(%s) failed: %v", assignedPod, err)
}
assigned.Status.State = ateapipb.WorkerState_WORKER_STATE_DRAINING
if err := tc.persistence.UpdateWorker(context.Background(), assigned, assigned.GetMetadata().GetVersion()); err != nil {
if _, err := tc.persistence.UpdateWorker(context.Background(), assigned.GetMetadata().GetName(), store.PreconditionFrom(assigned), func(toUpdate *ateapipb.Worker) error {
toUpdate.Status.State = ateapipb.WorkerState_WORKER_STATE_DRAINING
return nil
}); err != nil {
t.Fatalf("marking worker %s draining failed: %v", assignedPod, err)
}

Expand Down
4 changes: 4 additions & 0 deletions cmd/ateapi/internal/controlapi/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,10 @@ type serviceStore interface {
ListActorTemplates(ctx context.Context, atespace string, opts store.ListOptions) (store.ListResponse[*ateapipb.ActorTemplate], error)
DeleteActorTemplate(ctx context.Context, templateRef resources.ActorTemplateRef) (*ateapipb.ActorTemplate, error)
ListWorkers(ctx context.Context, opts store.ListOptions) (store.ListResponse[*ateapipb.Worker], error)
GetWorker(ctx context.Context, name string) (*ateapipb.Worker, error)
CreateWorker(ctx context.Context, worker *ateapipb.Worker) (*ateapipb.Worker, error)
UpdateWorker(ctx context.Context, name string, precondition store.Precondition, mutate func(toUpdate *ateapipb.Worker) error) (*ateapipb.Worker, error)
DeleteWorker(ctx context.Context, name string, pre store.DeletePreconditions) (*ateapipb.Worker, error)
AcquireLock(ctx context.Context, key string) (*store.Lock, error)
}

Expand Down
41 changes: 29 additions & 12 deletions cmd/ateapi/internal/controlapi/syncer.go
Original file line number Diff line number Diff line change
Expand Up @@ -91,9 +91,9 @@ type workerPoolSyncerStore interface {
GetActor(ctx context.Context, actorRef resources.ActorRef) (*ateapipb.Actor, error)
UpdateActor(ctx context.Context, actorRef resources.ActorRef, precondition store.Precondition, mutate func(toUpdate *ateapipb.Actor) error) (*ateapipb.Actor, error)
GetWorker(ctx context.Context, name string) (*ateapipb.Worker, error)
CreateWorker(ctx context.Context, worker *ateapipb.Worker) error
UpdateWorker(ctx context.Context, worker *ateapipb.Worker, expectedVersion int64) error
DeleteWorker(ctx context.Context, name string) error
CreateWorker(ctx context.Context, worker *ateapipb.Worker) (*ateapipb.Worker, error)
UpdateWorker(ctx context.Context, name string, precondition store.Precondition, mutate func(toUpdate *ateapipb.Worker) error) (*ateapipb.Worker, error)
DeleteWorker(ctx context.Context, name string, pre store.DeletePreconditions) (*ateapipb.Worker, error)
ListWorkers(ctx context.Context, opts store.ListOptions) (store.ListResponse[*ateapipb.Worker], error)
}

Expand Down Expand Up @@ -263,18 +263,19 @@ func (s *WorkerPoolSyncer) createOrUpdateWorker(ctx context.Context, key workerK
State: ateapipb.WorkerState_WORKER_STATE_ACTIVE,
},
}
// TODO(thockin): for now this is the only place Workers are
// created. If/when this becomes a regular API, validation should
// move there.
if errs := resources.ValidateWorker(worker, nil); len(errs) > 0 {
// TODO: validateWorker now lives next to CreateWorker, which applies it
// too. Once this path calls the RPC instead of the store, the check
// here goes away and the errors below arrive as INVALID_ARGUMENT.
if errs := validateWorker(worker, nil); len(errs) > 0 {
// Terminal: the inputs are deterministic, retrying cannot help. A
// future pod event re-enqueues the key.
slog.ErrorContext(ctx, "Invalid worker", append(key.logAttrs(), slog.Any("err", errs.ToAggregate()))...)
return nil
}
// ErrAlreadyExists means we lost a create race; requeue and converge
// via the update path.
return s.persistence.CreateWorker(ctx, worker)
_, err := s.persistence.CreateWorker(ctx, worker)
return err
}

changed := false
Expand All @@ -300,7 +301,13 @@ func (s *WorkerPoolSyncer) createOrUpdateWorker(ctx context.Context, key workerK

// ErrVersionConflict requeues the key; the retry re-fetches the worker at
// its new version.
return s.persistence.UpdateWorker(ctx, w, w.GetMetadata().GetVersion())
_, err = s.persistence.UpdateWorker(ctx, key.workerName(), store.PreconditionFrom(w), func(toUpdate *ateapipb.Worker) error {
toUpdate.Ip = w.GetIp()
toUpdate.SandboxClass = w.GetSandboxClass()
toUpdate.Labels = w.GetLabels()
return nil
})
return err
}

func isWorkerEligible(pod *corev1.Pod) bool {
Expand Down Expand Up @@ -356,8 +363,11 @@ func (s *WorkerPoolSyncer) markWorkerDraining(ctx context.Context, key workerKey
return nil
}
slog.InfoContext(ctx, "Syncer: marking worker draining (pod deleting)", key.logAttrs()...)
worker.Status.State = ateapipb.WorkerState_WORKER_STATE_DRAINING
return s.persistence.UpdateWorker(ctx, worker, worker.GetMetadata().GetVersion())
_, err = s.persistence.UpdateWorker(ctx, key.workerName(), store.PreconditionFrom(worker), func(toUpdate *ateapipb.Worker) error {
toUpdate.Status.State = ateapipb.WorkerState_WORKER_STATE_DRAINING
return nil
})
return err
}

// reconcileDeadWorker cleans up a worker whose pod is gone. It releases the
Expand All @@ -370,7 +380,14 @@ func (s *WorkerPoolSyncer) reconcileDeadWorker(ctx context.Context, name string)
if err := s.releaseActorOnDeadWorker(ctx, name); err != nil {
return err
}
return s.persistence.DeleteWorker(ctx, name)
// The delete now reports absence rather than succeeding silently, but a
// worker already gone is exactly the state this is driving towards.
// Idempotency lives here, at the caller, so re-driving a reconcile is safe.
_, err := s.persistence.DeleteWorker(ctx, name, store.DeletePreconditions{})
if errors.Is(err, store.ErrNotFound) {
return nil
}
return err
}

// storedWorkerListBackoff and storedWorkerListCap are the exponential backoff
Expand Down
Loading
Loading