mirror of
https://github.com/agent-substrate/substrate.git
synced 2026-10-02 03:24:42 +08:00
fix(ateapi): add store fallback in ActorIdentity to resolve test flakiness (#1113)
Fix `TestDurableDirLifecycle` flaky test https://github.com/agent-substrate/substrate/actions/runs/32416349094/job/96578260374 Cause: - In the specific test (`_suspend_from_PAUSED`), the actor is suspended (which strips its worker assignment, returning the worker to the pool) and then immediately resumed. - When the test calls `ResumeActor`, the control plane selects a new worker, writes the assignment to the Kubernetes database, and synchronously makes a direct gRPC Restore() call to the atelet on that new worker. - The `atelet` boots the actor and `atunnel`. `atunnel` then immediately asks the `ateapi` (via the `ActorIdentity` service) to mint an actor certificate so it can connect to the network. Race: The `ActorIdentity` service verifies that the worker is actually assigned to the actor by checking its internal `workercache.Cache`. Because `ResumeActor` was so fast, the cache hasn't processed the new worker assignment yet. The `ActorIdentity` service looks at the stale cache, sees the worker is unassigned, and hard-rejects the certificate minting with caller is not permitted to mint credentials for this actor. I faced this issue too when working on the postgresql updateWorker optimization (#934) Fix: Updated `authorizeActor` inside `cmd/ateapi/internal/actoridentity/actoridentity.go`. Now, if `authorizeActor` fails an authorization check (like discovering a missing or mismatched worker assignment), it will automatically do a one-time "live fetch" bypass (s.store.GetWorker()) directly against the primary store to bypass the lagging informer cache. - [ ] Tests pass - [ ] Appropriate changes to documentation are included in the PR
This commit is contained in:
@@ -317,38 +317,49 @@ func validateWorkerRef(worker *ateapipb.ObjectRef) error {
|
||||
// authorizeActor resolves the actor from the authenticated worker and verifies
|
||||
// that the worker and actor still point at one another. Actor identity supplied
|
||||
// by the requester never participates in this authorization decision.
|
||||
// The worker is resolved from cache first (hot path), but denials fall back
|
||||
// to the authoritative store to handle watch-delivery lag right after ResumeActor.
|
||||
// The worker is resolved from cache first (hot path), but cache misses and
|
||||
// denials fall back to the authoritative store to handle watch-delivery lag
|
||||
// right after ResumeActor.
|
||||
func (s *Server) authorizeActor(ctx context.Context, caller *ateletCaller, req *ateapipb.MintCertRequest) (*ateapipb.Actor, resources.ActorRef, error) {
|
||||
reason := "worker not found"
|
||||
worker, err := s.workers.Worker(req.GetWorker().GetName())
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
// Cache misses cannot read through without the pool name; worker rows exist for the pod's lifetime.
|
||||
return nil, resources.ActorRef{}, s.denyMint(ctx, caller, req, "worker not found")
|
||||
}
|
||||
if err != nil && !errors.Is(err, store.ErrNotFound) {
|
||||
slog.ErrorContext(ctx, "ActorIdentity: failed to read worker", slog.Any("err", err))
|
||||
return nil, resources.ActorRef{}, status.Error(codes.Internal, "failed to look up worker")
|
||||
}
|
||||
|
||||
actor, actorRef, err := s.authorizeWithWorker(ctx, worker, caller, req)
|
||||
if status.Code(err) != codes.PermissionDenied {
|
||||
return actor, actorRef, err // authorized, or a non-denial error
|
||||
if err == nil {
|
||||
actor, actorRef, mismatchReason, err := s.authorizeWithWorker(ctx, worker, caller, req)
|
||||
if err == nil {
|
||||
return actor, actorRef, nil
|
||||
}
|
||||
if !errors.Is(err, errAssignmentMismatch) {
|
||||
return nil, resources.ActorRef{}, err // e.g. actor lookup failed
|
||||
}
|
||||
reason = mismatchReason
|
||||
}
|
||||
|
||||
// Read-through: re-check the authoritative worker from the store on denial.
|
||||
// Read-through: re-check the authoritative worker from the store on a
|
||||
// cache miss or assignment mismatch. Only fresh data may authorize, and
|
||||
// only fresh data may deny.
|
||||
fresh, ferr := s.store.GetWorker(ctx, req.GetWorker().GetName())
|
||||
if ferr != nil {
|
||||
if !errors.Is(ferr, store.ErrNotFound) {
|
||||
slog.ErrorContext(ctx, "ActorIdentity: read-through worker lookup failed", slog.Any("err", ferr))
|
||||
}
|
||||
return nil, resources.ActorRef{}, err // the original denial stands
|
||||
return nil, resources.ActorRef{}, s.denyMint(ctx, caller, req, reason) // the cached verdict stands
|
||||
}
|
||||
actor, actorRef, retryErr := s.authorizeWithWorker(ctx, fresh, caller, req)
|
||||
if retryErr == nil {
|
||||
slog.InfoContext(ctx, "ActorIdentity: authorized via store read-through; worker cache was stale",
|
||||
slog.String("worker", req.GetWorker().GetName()))
|
||||
|
||||
actor, actorRef, retryReason, retryErr := s.authorizeWithWorker(ctx, fresh, caller, req)
|
||||
if retryErr != nil {
|
||||
if errors.Is(retryErr, errAssignmentMismatch) {
|
||||
return nil, resources.ActorRef{}, s.denyMint(ctx, caller, req, retryReason)
|
||||
}
|
||||
return nil, resources.ActorRef{}, retryErr
|
||||
}
|
||||
return actor, actorRef, retryErr
|
||||
|
||||
slog.InfoContext(ctx, "ActorIdentity: authorized via store read-through; worker cache was stale",
|
||||
slog.String("worker", req.GetWorker().GetName()))
|
||||
return actor, actorRef, nil
|
||||
}
|
||||
|
||||
// denyMint logs the internal reason and returns a uniform PermissionDenied.
|
||||
@@ -360,49 +371,46 @@ func (s *Server) denyMint(ctx context.Context, caller *ateletCaller, req *ateapi
|
||||
return status.Errorf(codes.PermissionDenied, "caller is not permitted to mint credentials for this actor: %s", reason)
|
||||
}
|
||||
|
||||
// authorizeWithWorker runs the worker↔actor mutual-pointer checks against one
|
||||
// specific worker record (cached or freshly read — the caller decides).
|
||||
func (s *Server) authorizeWithWorker(ctx context.Context, worker *ateapipb.Worker, caller *ateletCaller, req *ateapipb.MintCertRequest) (*ateapipb.Actor, resources.ActorRef, error) {
|
||||
deny := func(reason string, args ...any) error {
|
||||
return s.denyMint(ctx, caller, req, reason, args...)
|
||||
}
|
||||
var errAssignmentMismatch = errors.New("assignment mismatch")
|
||||
|
||||
// authorizeWithWorker returns errAssignmentMismatch and a reason string if the authorization failed
|
||||
// due to an assignment mismatch, indicating the caller may want to refetch the worker and retry.
|
||||
func (s *Server) authorizeWithWorker(ctx context.Context, worker *ateapipb.Worker, caller *ateletCaller, req *ateapipb.MintCertRequest) (*ateapipb.Actor, resources.ActorRef, string, error) {
|
||||
if worker.GetNodeName() != caller.nodeName {
|
||||
return nil, resources.ActorRef{}, deny("worker is hosted on a different node", slog.String("workerNode", worker.GetNodeName()))
|
||||
return nil, resources.ActorRef{}, "worker is hosted on a different node", errAssignmentMismatch
|
||||
}
|
||||
|
||||
actorRef := resources.ActorRefFromObjectRef(worker.GetStatus().GetAssignment().GetActor())
|
||||
if actorRef == (resources.ActorRef{}) {
|
||||
return nil, resources.ActorRef{}, deny("worker has no actor assignment")
|
||||
return nil, resources.ActorRef{}, "worker has no actor assignment", errAssignmentMismatch
|
||||
}
|
||||
|
||||
actor, err := s.store.GetActor(ctx, actorRef)
|
||||
if err != nil {
|
||||
if errors.Is(err, store.ErrNotFound) {
|
||||
return nil, resources.ActorRef{}, deny("assigned actor not found")
|
||||
return nil, resources.ActorRef{}, "assigned actor not found", errAssignmentMismatch
|
||||
}
|
||||
slog.ErrorContext(ctx, "ActorIdentity: failed to read actor", slog.Any("actor", actorRef), slog.Any("err", err))
|
||||
return nil, resources.ActorRef{}, status.Error(codes.Internal, "failed to look up actor")
|
||||
return nil, resources.ActorRef{}, "", status.Error(codes.Internal, "failed to look up actor")
|
||||
}
|
||||
|
||||
// Refuse credential minting if the actor is being deleted. Under force deletion,
|
||||
// an actor enters ACTOR_STATE_DELETING while its worker assignment is still active.
|
||||
if actor.GetStatus().GetState() == ateapipb.ActorState_ACTOR_STATE_DELETING {
|
||||
slog.WarnContext(ctx, "ActorIdentity refused: actor is being deleted", slog.Any("actor", actorRef))
|
||||
return nil, resources.ActorRef{}, status.Error(codes.FailedPrecondition, "actor is being deleted")
|
||||
return nil, resources.ActorRef{}, "", status.Error(codes.FailedPrecondition, "actor is being deleted")
|
||||
}
|
||||
|
||||
// An actor placed on a worker always carries its placement fields. Missing
|
||||
// placement is a control-plane bug rather than a client error, so it is not
|
||||
// folded into deny().
|
||||
assignment := actor.GetStatus().GetWorkerAssignment()
|
||||
if assignment == nil {
|
||||
slog.ErrorContext(ctx, "ActorIdentity: running actor has no worker assignment", slog.Any("actor", actorRef))
|
||||
return nil, resources.ActorRef{}, status.Error(codes.FailedPrecondition, "actor has no worker assigned")
|
||||
return nil, resources.ActorRef{}, "", status.Error(codes.FailedPrecondition, "actor has no worker assigned")
|
||||
}
|
||||
if worker.GetStatus().GetAssignment().GetActorUid() != actor.GetMetadata().GetUid() {
|
||||
return nil, resources.ActorRef{}, deny("worker is no longer assigned to this actor incarnation", slog.Any("actor", actorRef))
|
||||
return nil, resources.ActorRef{}, "worker is no longer assigned to this actor incarnation", errAssignmentMismatch
|
||||
}
|
||||
if assignment.GetWorker().GetName() != worker.GetMetadata().GetName() {
|
||||
return nil, resources.ActorRef{}, deny("actor no longer points to the requesting worker", slog.Any("actor", actorRef))
|
||||
return nil, resources.ActorRef{}, "actor no longer points to the requesting worker", errAssignmentMismatch
|
||||
}
|
||||
return actor, actorRef, nil
|
||||
return actor, actorRef, "", nil
|
||||
}
|
||||
|
||||
@@ -241,6 +241,74 @@ func TestMintCertReadsThroughStaleWorkerCache(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestMintCertReadsThroughWorkerCacheMiss pins the read-through for a worker
|
||||
// the cache has never seen: a worker registered moments before assignment may
|
||||
// be committed to the store (possibly by another replica) before this
|
||||
// replica's cache has received the worker row at all. Absence from the cache
|
||||
// is stale data and must not deny by itself; absence from the store must.
|
||||
func TestMintCertReadsThroughWorkerCacheMiss(t *testing.T) {
|
||||
for name, workerInStore := range map[string]bool{
|
||||
"worker assigned in store but not yet in cache: authorized via read-through": true,
|
||||
"worker in neither cache nor store: denial stands": false,
|
||||
} {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
st, cleanup := storetest.SetupTestStore(t)
|
||||
defer cleanup()
|
||||
|
||||
// Phase 1: only the actor exists; the cache seeds with no workers
|
||||
// and (via the inert watch) never learns of any.
|
||||
seedActor(t, ctx, st, actorFixture{state: ateapipb.ActorState_ACTOR_STATE_RUNNING, workerNode: testNode, noWorker: true})
|
||||
workers := workercache.New(staleWatchStore{st}, time.Hour)
|
||||
cacheCtx, cancel := context.WithCancel(ctx)
|
||||
t.Cleanup(cancel)
|
||||
if err := workers.Start(cacheCtx); err != nil {
|
||||
t.Fatalf("start worker cache: %v", err)
|
||||
}
|
||||
|
||||
actor, err := st.GetActor(ctx, resources.ActorRef{Atespace: testAtespace, Name: testActorName})
|
||||
if err != nil {
|
||||
t.Fatalf("read seeded actor: %v", err)
|
||||
}
|
||||
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{
|
||||
Metadata: &ateapipb.ResourceMetadata{Name: testWorkerName},
|
||||
WorkerNamespace: testPodNS,
|
||||
WorkerPool: testPool,
|
||||
WorkerPod: testWorkerPod,
|
||||
WorkerPodUid: testWorkerPodUID,
|
||||
NodeName: testNode,
|
||||
Status: &ateapipb.WorkerStatus{
|
||||
State: ateapipb.WorkerState_WORKER_STATE_ACTIVE,
|
||||
Assignment: &ateapipb.ActorAssignment{
|
||||
Actor: (resources.ActorRef{Atespace: testAtespace, Name: testActorName}).ToObjectRef(),
|
||||
ActorUid: actor.GetMetadata().GetUid(),
|
||||
},
|
||||
},
|
||||
}); err != nil {
|
||||
t.Fatalf("register worker in store: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
srv := newTestServerWithCache(t, st, workers)
|
||||
resp, err := srv.MintCert(ctxWithCert(ateletCertOn(t, testNode)), mintCertRequest(t, actor.GetMetadata().GetUid()))
|
||||
|
||||
wantCode := codes.PermissionDenied
|
||||
if workerInStore {
|
||||
wantCode = codes.OK
|
||||
}
|
||||
if got := status.Code(err); got != wantCode {
|
||||
t.Fatalf("MintCert() code = %v (err = %v), want %v", got, err, wantCode)
|
||||
}
|
||||
if workerInStore && len(resp.GetActorCertificates()) == 0 {
|
||||
t.Fatal("MintCert() returned no certificates")
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// newTestServerWithCache is newTestServer with a caller-controlled worker
|
||||
// cache (e.g. one frozen at a stale state).
|
||||
func newTestServerWithCache(t *testing.T, st store.Interface, workers *workercache.Cache) *Server {
|
||||
|
||||
Reference in New Issue
Block a user