mirror of
https://github.com/mudler/LocalAI.git
synced 2026-09-22 06:04:55 -04:00
GET /api/cluster/peer authenticated with the deployment's shared registration token and took the dialling replica's id from ?id= on trust. Every worker holds that token, so anything holding it could open a peer link as any replica: relay through it to every worker tunnel that replica owns, displace a real replica's inbound link by declaring its id, and point the roughly 31 GiB per-session receive window at one replica. Validating the id against the instances table does not fix this, because the attack declares a real replica's id. So the route now checks two credentials and needs both. The shared token still says the dialler belongs to this deployment; a new per-replica credential says which replica it is. The credential follows the per-node worker credential rather than inventing a second mechanism: crypto/rand.Text, stored only as a hex SHA-256, compared in constant time, with no fallback to the shared token. It differs in the stronger direction. A worker's credential is minted by the frontend and handed over once; a replica writes its own instances row, so it mints its own secret, publishes only the hash in the same statement that publishes its address, and never sends the plaintext anywhere but the peer dial. A peer that presents no credential is refused, not waved through. An old replica and an attacker holding the shared token send the same request, so accepting the first accepts the second; there is no safe downgrade here, only a quiet one. The refusal is made loud instead, on both sides, naming the upgrade rather than the network. On the documented frontend-first order a new replica still dials an old one; an old replica cannot dial a new one, which costs relayed requests that land on a not-yet-restarted replica and surfaces as no route, never as absence. A rejected peer gets its own sentinel, ErrPeerRejected, whose unwrap chain carries ErrPeerUnreachable as well and no absence sentinel at all. Keeping the older sentinel means no existing consumer changes behaviour; the cause stays out of the chain, so absence cannot escape through it and nothing can read an authorization failure as a worker that went away. One consequence beyond the fix: a replica with no advertised address has no instances row, so it now cannot dial out either. It was already unreachable inward. The startup error and the docs say so. Registry.Register, NewMembership, NewPeerPool, PeerHandler and RegisterClusterRoutes all gained required arguments, so the identity cannot be dropped without a compile failure. Assisted-by: Claude Opus 5 [claude-code] Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
472 lines
21 KiB
Go
472 lines
21 KiB
Go
package distributed_test
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"io"
|
|
"net"
|
|
"net/http"
|
|
"net/url"
|
|
"strings"
|
|
"time"
|
|
|
|
clustersvc "github.com/mudler/LocalAI/core/services/cluster"
|
|
|
|
"github.com/gorilla/websocket"
|
|
"github.com/libp2p/go-yamux/v5"
|
|
|
|
. "github.com/onsi/ginkgo/v2"
|
|
. "github.com/onsi/gomega"
|
|
"gorm.io/driver/postgres"
|
|
"gorm.io/gorm"
|
|
gormlogger "gorm.io/gorm/logger"
|
|
)
|
|
|
|
const (
|
|
// instanceRosterTimeout bounds the wait for a replica's row to appear.
|
|
// Registration is synchronous in startup, so this only has to cover the gap
|
|
// between /readyz answering and this spec's first query.
|
|
instanceRosterTimeout = "30s"
|
|
instanceRosterPoll = "500ms"
|
|
|
|
// deadReplicaTimeout bounds the wait for a survivor to reap a replica that
|
|
// was killed: the liveness window plus a sweep interval plus slack. It is
|
|
// deliberately derived from the constants rather than a round number, so
|
|
// tightening the window shortens the spec instead of leaving it passing for
|
|
// the wrong reason.
|
|
deadReplicaTimeout = clustersvc.InstanceLiveness + 4*clustersvc.InstanceHeartbeat
|
|
|
|
// peerDialTimeout bounds one peer dial. Every replica here is a local
|
|
// process, so a dial that needs longer has failed, not slowed.
|
|
peerDialTimeout = 20 * time.Second
|
|
|
|
// gracefulDepartureTimeout bounds the wait for a cleanly stopped replica to
|
|
// leave the table. It must stay well under InstanceLiveness, which the spec
|
|
// asserts: a budget that reached the window would pass on the sweeper doing
|
|
// the work and prove nothing about deregistration.
|
|
gracefulDepartureTimeout = 15 * time.Second
|
|
|
|
// peerRefusalTimeout bounds how long a refused stream may take to end. It
|
|
// is short on purpose: the refusal is one frame from a replica that has
|
|
// already decided, so a stream still open at this point is parked.
|
|
peerRefusalTimeout = 5 * time.Second
|
|
|
|
// unheldNodeID is a worker id no replica holds a tunnel for. It is a
|
|
// well-formed id rather than a nonsense string so the refusal it draws is
|
|
// the routing answer and not a parse failure.
|
|
unheldNodeID = "00000000-0000-0000-0000-00000000dead"
|
|
)
|
|
|
|
// openClusterDB connects to the database the cluster was given, so a spec can
|
|
// read the tables the peer link keeps. Nothing serves them over HTTP: they are
|
|
// replica-to-replica state, not an admin surface, and inventing an endpoint to
|
|
// observe them would be a bigger change than the thing under test.
|
|
func openClusterDB(dsn string) *gorm.DB {
|
|
GinkgoHelper()
|
|
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{Logger: gormlogger.Discard})
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(func() { closeDB(db) })
|
|
return db
|
|
}
|
|
|
|
// hostPortOf strips the scheme off a frontend URL, giving the form the
|
|
// instances table stores.
|
|
func hostPortOf(url string) string {
|
|
return strings.TrimPrefix(strings.TrimPrefix(url, "http://"), "https://")
|
|
}
|
|
|
|
// instanceRoster reads the live replica rows, keeping the last error so a
|
|
// failing Eventually can name it.
|
|
type instanceRoster struct {
|
|
registry *clustersvc.Registry
|
|
ctx context.Context
|
|
|
|
lastErr error
|
|
lastSaw []clustersvc.Instance
|
|
}
|
|
|
|
func newInstanceRoster(db *gorm.DB) *instanceRoster {
|
|
return &instanceRoster{registry: clustersvc.NewRegistry(db), ctx: context.Background()}
|
|
}
|
|
|
|
// addresses returns the advertised address of every live replica, or nil on a
|
|
// query error so Eventually keeps trying.
|
|
func (r *instanceRoster) addresses() []string {
|
|
live, err := r.registry.Live(r.ctx, clustersvc.InstanceLiveness)
|
|
if err != nil {
|
|
r.lastErr = err
|
|
return nil
|
|
}
|
|
r.lastErr = nil
|
|
r.lastSaw = live
|
|
addrs := []string{}
|
|
for _, instance := range live {
|
|
addrs = append(addrs, instance.AdvertisedAddr)
|
|
}
|
|
return addrs
|
|
}
|
|
|
|
// idAt returns the id of the live replica advertising addr, or "" if no such
|
|
// row is present yet.
|
|
func (r *instanceRoster) idAt(addr string) string {
|
|
for _, instance := range r.lastSaw {
|
|
if instance.AdvertisedAddr == addr {
|
|
return instance.ID
|
|
}
|
|
}
|
|
return ""
|
|
}
|
|
|
|
func (r *instanceRoster) describe() string {
|
|
if r.lastErr != nil {
|
|
return fmt.Sprintf("the last read of the instances table failed: %v", r.lastErr)
|
|
}
|
|
return fmt.Sprintf("the instances table held %d live replica(s): %+v", len(r.lastSaw), r.lastSaw)
|
|
}
|
|
|
|
// awaitReplicas waits for every frontend of c to publish its address and
|
|
// returns the roster, positioned on that reading.
|
|
func awaitReplicas(roster *instanceRoster, addrs ...string) {
|
|
GinkgoHelper()
|
|
Eventually(roster.addresses, instanceRosterTimeout, instanceRosterPoll).
|
|
Should(ConsistOf(addrs), roster.describe)
|
|
}
|
|
|
|
// joinAsPeer gives this spec process a replica identity of its own: a freshly
|
|
// minted credential, with only its hash published in the instances table.
|
|
//
|
|
// It exists because the peer route no longer takes ?id= on trust. A spec that
|
|
// plays a sibling replica has to join the cluster the way a replica does, which
|
|
// is the point rather than a chore: the address it publishes is never dialled,
|
|
// but the credential it publishes is what every dial below is checked against.
|
|
func joinAsPeer(roster *instanceRoster, id string) clustersvc.PeerCredential {
|
|
GinkgoHelper()
|
|
cred := clustersvc.NewPeerCredential()
|
|
Expect(roster.registry.Register(roster.ctx, id, "127.0.0.1:1", "e2e", cred.Hash())).To(Succeed())
|
|
return cred
|
|
}
|
|
|
|
var _ = Describe("Cluster peer link", Label("Distributed"), Label("Cluster"), func() {
|
|
It("publishes an address for every replica that peers can actually dial", func() {
|
|
// A wrong implementation registers nothing (the whole of phase 1 had no
|
|
// call site until this spec), registers one row for two replicas, or
|
|
// records an address nothing can connect to: the bind address of a
|
|
// replica behind a service, or the loopback address the route to a
|
|
// co-located database would suggest.
|
|
c, dsn := startClusterOnFreshDB(2, 0)
|
|
|
|
roster := newInstanceRoster(openClusterDB(dsn))
|
|
awaitReplicas(roster, hostPortOf(c.FrontendURL(0)), hostPortOf(c.FrontendURL(1)))
|
|
|
|
// "Routable" is not a property of the string. Connect to each address,
|
|
// which is the only check that would have caught a replica publishing
|
|
// the port it was configured with rather than the one it serves on.
|
|
for _, instance := range roster.lastSaw {
|
|
conn, err := net.DialTimeout("tcp", instance.AdvertisedAddr, peerDialTimeout)
|
|
Expect(err).ToNot(HaveOccurred(),
|
|
"replica %s advertises %q, which nothing can connect to", instance.ID, instance.AdvertisedAddr)
|
|
Expect(conn.Close()).To(Succeed())
|
|
}
|
|
})
|
|
|
|
It("carries a peer stream between two replicas, and refuses one without the cluster token", func() {
|
|
// A wrong implementation fails here on WebSocket framing, which is the
|
|
// likeliest defect in the peer link: the adapter has to turn
|
|
// message-oriented WebSocket frames into the undelimited byte stream
|
|
// yamux drives. It also fails if the route was never registered on the
|
|
// real server, or if the global session middleware answers it: a peer
|
|
// carries no session and no user, only the cluster token.
|
|
//
|
|
// The stream is opened with the production dialler, resolving the peer
|
|
// through the production registry, over a real socket to a real
|
|
// process. This spec plays the sibling replica, because phase 1 has
|
|
// nothing that makes a frontend dial one on its own.
|
|
c, dsn := startClusterOnFreshDB(2, 0)
|
|
|
|
roster := newInstanceRoster(openClusterDB(dsn))
|
|
awaitReplicas(roster, hostPortOf(c.FrontendURL(0)), hostPortOf(c.FrontendURL(1)))
|
|
|
|
peerID := roster.idAt(hostPortOf(c.FrontendURL(1)))
|
|
Expect(peerID).ToNot(BeEmpty())
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), peerDialTimeout)
|
|
defer cancel()
|
|
|
|
// The spec joins as a replica before it dials as one. Without this the
|
|
// frontend refuses the link: the id would be a claim it cannot check,
|
|
// which is exactly what a peer dial is now required to prove.
|
|
selfCred := joinAsPeer(roster, "e2e-peer")
|
|
pool := clustersvc.NewPeerPool("e2e-peer", c.RegistrationToken(), selfCred, roster.registry)
|
|
DeferCleanup(pool.Close)
|
|
|
|
// OpenStream is only acknowledged once the far side accepts, so this
|
|
// returning at all proves the frontend is accepting streams on the
|
|
// session it took, in addition to proving the handshake.
|
|
stream, err := pool.Open(ctx, peerID)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(func() { _ = stream.Close() })
|
|
|
|
// Phase 2 installs the relay on this link, so an accepted stream is one
|
|
// the peer is waiting to be told which worker it is for. Name one no
|
|
// replica holds and the refusal must come back at once.
|
|
//
|
|
// This spec used to assert the opposite, that an accepted stream ended
|
|
// immediately, because phase 1 had no relay to hand it to. The relay
|
|
// made that stale rather than wrong: a stream that says nothing now
|
|
// parks for relayHeaderTimeout, which is 15 seconds, and the old
|
|
// assertion failed on a five second budget against a replica behaving
|
|
// exactly as designed.
|
|
Expect(stream.SetWriteDeadline(time.Now().Add(peerRefusalTimeout))).To(Succeed())
|
|
Expect(clustersvc.WriteRelayRequest(stream, unheldNodeID, peerRefusalTimeout)).To(Succeed())
|
|
|
|
Expect(stream.SetReadDeadline(time.Now().Add(peerRefusalTimeout))).To(Succeed())
|
|
err = clustersvc.ReadRelayReply(stream)
|
|
Expect(err).To(MatchError(clustersvc.ErrNotOwner),
|
|
"the peer did not refuse a worker it does not hold: %v", err)
|
|
Expect(err).ToNot(MatchError(clustersvc.ErrNoConnection),
|
|
"a replica that does not hold a tunnel must not report the worker as absent: that is how a scheduler evicts a healthy worker")
|
|
|
|
// And the refusal ENDS the stream. A replica that says why and leaves
|
|
// the stream open has parked the caller on a request that will never be
|
|
// served, which reads as a slow replica rather than a refused request,
|
|
// and no deadline on the far side can tell those apart.
|
|
Expect(stream.SetReadDeadline(time.Now().Add(peerRefusalTimeout))).To(Succeed())
|
|
_, err = stream.Read(make([]byte, 1))
|
|
Expect(err).To(SatisfyAny(MatchError(io.EOF), MatchError(yamux.ErrStreamReset)),
|
|
"the peer refused the stream and then left it open: %v", err)
|
|
|
|
// The same dial with the wrong shared token must be refused, otherwise
|
|
// the success above says nothing about authentication. The identity
|
|
// credential is correct here, so what this pins is that adding identity
|
|
// did not replace the token check with it.
|
|
impostor := clustersvc.NewPeerPool("e2e-peer", "not-the-cluster-token", selfCred, roster.registry)
|
|
DeferCleanup(impostor.Close)
|
|
_, err = impostor.Open(ctx, peerID)
|
|
Expect(err).To(MatchError(clustersvc.ErrPeerUnreachable))
|
|
Expect(err).ToNot(MatchError(clustersvc.ErrInstanceNotFound),
|
|
"a peer refusing credentials is a live peer; reading it as absence is how a replica evicts healthy workers")
|
|
})
|
|
|
|
It("refuses a peer that declares a replica id it cannot prove, and leaves that replica's link alone", func() {
|
|
// The security gap this route carried for the whole phase, closed here
|
|
// against real processes. The attacker's position is the realistic one:
|
|
// it holds LOCALAI_REGISTRATION_TOKEN, which every worker in the
|
|
// deployment holds, and it knows a replica id, which any log line or
|
|
// roster read gives it. Before per-replica credentials that was the
|
|
// entire check, so it could open a link and relay through it to every
|
|
// worker tunnel the replica it dialled owns, and by declaring somebody
|
|
// else's id evict that replica's inbound link at will.
|
|
//
|
|
// It is also the case the cheap fix does not reach. Validating ?id=
|
|
// against the instances table passes every dial below except the last:
|
|
// the ids ARE in the table. Only a secret the impostor does not hold
|
|
// refuses them.
|
|
c, dsn := startClusterOnFreshDB(2, 0)
|
|
|
|
roster := newInstanceRoster(openClusterDB(dsn))
|
|
awaitReplicas(roster, hostPortOf(c.FrontendURL(0)), hostPortOf(c.FrontendURL(1)))
|
|
|
|
target := roster.idAt(hostPortOf(c.FrontendURL(1)))
|
|
victim := roster.idAt(hostPortOf(c.FrontendURL(0)))
|
|
Expect(target).ToNot(BeEmpty())
|
|
Expect(victim).ToNot(BeEmpty())
|
|
|
|
ctx, cancel := context.WithTimeout(context.Background(), peerDialTimeout)
|
|
defer cancel()
|
|
|
|
// A legitimate link first, so there is something for an impostor to
|
|
// displace and so the refusals below are not passing because the route
|
|
// is broken for everyone.
|
|
selfCred := joinAsPeer(roster, "e2e-peer")
|
|
legit := clustersvc.NewPeerPool("e2e-peer", c.RegistrationToken(), selfCred, roster.registry)
|
|
DeferCleanup(legit.Close)
|
|
held, err := legit.Open(ctx, target)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
DeferCleanup(func() { _ = held.Close() })
|
|
|
|
// Declaring a REAL frontend replica's id. This is the dial that used to
|
|
// succeed, and the one a scheduler must never be told anything about a
|
|
// worker from.
|
|
asReplica := clustersvc.NewPeerPool(victim, c.RegistrationToken(), clustersvc.NewPeerCredential(), roster.registry)
|
|
DeferCleanup(asReplica.Close)
|
|
_, err = asReplica.Open(ctx, target)
|
|
Expect(err).To(MatchError(clustersvc.ErrPeerRejected),
|
|
"a frontend accepted a peer link from something that could not prove it was the replica it named")
|
|
Expect(err).ToNot(MatchError(clustersvc.ErrInstanceNotFound),
|
|
"a rejected peer is an authorization failure; read as absence it makes the scheduler reap rows and evict models")
|
|
Expect(err).ToNot(MatchError(clustersvc.ErrNoConnection),
|
|
"a rejected peer must never be readable as a worker with no connection")
|
|
|
|
// And declaring the id whose link is currently held, which is the
|
|
// eviction half: SessionStore keeps one link per id, so an accepted
|
|
// dial closes whatever was there.
|
|
displacer := clustersvc.NewPeerPool("e2e-peer", c.RegistrationToken(), clustersvc.NewPeerCredential(), roster.registry)
|
|
DeferCleanup(displacer.Close)
|
|
_, err = displacer.Open(ctx, target)
|
|
Expect(err).To(MatchError(clustersvc.ErrPeerRejected))
|
|
|
|
// The link opened before the attempts still carries a request. It is
|
|
// the stream taken above, not a fresh one: a fresh Open would simply
|
|
// re-dial and succeed, which would say nothing about whether the old
|
|
// session survived.
|
|
Expect(held.SetWriteDeadline(time.Now().Add(peerRefusalTimeout))).To(Succeed())
|
|
Expect(clustersvc.WriteRelayRequest(held, unheldNodeID, peerRefusalTimeout)).To(Succeed(),
|
|
"the impostor closed the real peer's link")
|
|
Expect(held.SetReadDeadline(time.Now().Add(peerRefusalTimeout))).To(Succeed())
|
|
Expect(clustersvc.ReadRelayReply(held)).To(MatchError(clustersvc.ErrNotOwner),
|
|
"the impostor closed the real peer's link")
|
|
})
|
|
|
|
It("refuses a peer that presents no credential at all, the way a replica running an older release dials", func() {
|
|
// The mixed-version window, made explicit. An upgraded frontend refuses
|
|
// a peer that sends no credential, rather than falling back to the
|
|
// shared token, and the refusal is a plain 401 before the WebSocket
|
|
// upgrade so the dialler reads a status and not a framing error.
|
|
//
|
|
// This is dialled raw rather than through PeerPool, because PeerPool on
|
|
// this release always sends the header: an old replica is the only
|
|
// thing that does not, and it is what this reproduces.
|
|
c, dsn := startClusterOnFreshDB(1, 0)
|
|
|
|
roster := newInstanceRoster(openClusterDB(dsn))
|
|
awaitReplicas(roster, hostPortOf(c.FrontendURL(0)))
|
|
target := roster.idAt(hostPortOf(c.FrontendURL(0)))
|
|
Expect(target).ToNot(BeEmpty())
|
|
|
|
// A credential that would be perfectly valid if it were sent.
|
|
joinAsPeer(roster, "e2e-old-peer")
|
|
|
|
endpoint := url.URL{
|
|
Scheme: "ws",
|
|
Host: hostPortOf(c.FrontendURL(0)),
|
|
Path: clustersvc.PeerPath,
|
|
RawQuery: url.Values{"id": []string{"e2e-old-peer"}}.Encode(),
|
|
}
|
|
header := http.Header{}
|
|
header.Set("Authorization", "Bearer "+c.RegistrationToken())
|
|
|
|
dialer := &websocket.Dialer{HandshakeTimeout: peerDialTimeout}
|
|
conn, resp, err := dialer.Dial(endpoint.String(), header)
|
|
if conn != nil {
|
|
DeferCleanup(func() { _ = conn.Close() })
|
|
}
|
|
Expect(resp).ToNot(BeNil())
|
|
if resp.Body != nil {
|
|
DeferCleanup(func() { _ = resp.Body.Close() })
|
|
}
|
|
Expect(err).To(HaveOccurred(),
|
|
"an upgraded frontend accepted a peer link carrying only the shared token")
|
|
Expect(resp.StatusCode).To(Equal(http.StatusUnauthorized))
|
|
Expect(resp.Header.Get("Upgrade")).To(BeEmpty(),
|
|
"the frontend upgraded a dial it then refused; a peer reads the status, not a WebSocket error")
|
|
})
|
|
|
|
It("stops being dialled as soon as a replica shuts down cleanly", func() {
|
|
// The crash case below is handled by the sweeper, at the cost of a
|
|
// whole liveness window of peers dialling a corpse. A rolling update is
|
|
// not a crash: the replica knows it is leaving and says so. Without
|
|
// deregistration the two are indistinguishable, and every rolling
|
|
// restart spends that window failing peer dials for no reason.
|
|
c, dsn := startClusterOnFreshDB(2, 0)
|
|
|
|
roster := newInstanceRoster(openClusterDB(dsn))
|
|
awaitReplicas(roster, hostPortOf(c.FrontendURL(0)), hostPortOf(c.FrontendURL(1)))
|
|
departingID := roster.idAt(hostPortOf(c.FrontendURL(1)))
|
|
Expect(departingID).ToNot(BeEmpty())
|
|
|
|
Expect(c.StopFrontendGracefully(1)).To(Succeed())
|
|
Eventually(func() bool { return c.FrontendAlive(1) }, "20s", "500ms").Should(BeFalse())
|
|
|
|
// The budget is deliberately shorter than the liveness window: passing
|
|
// it proves the replica announced its departure rather than aged out.
|
|
Expect(gracefulDepartureTimeout).To(BeNumerically("<", clustersvc.InstanceLiveness))
|
|
Eventually(roster.addresses, gracefulDepartureTimeout, instanceRosterPoll).
|
|
Should(ConsistOf(hostPortOf(c.FrontendURL(0))), roster.describe)
|
|
|
|
// And absence is the RIGHT answer here, unlike the killed case: the
|
|
// replica said it was going. A caller may act on this.
|
|
ctx, cancel := context.WithTimeout(context.Background(), peerDialTimeout)
|
|
defer cancel()
|
|
pool := clustersvc.NewPeerPool("e2e-peer", c.RegistrationToken(), joinAsPeer(roster, "e2e-peer"), roster.registry)
|
|
DeferCleanup(pool.Close)
|
|
_, err := pool.Open(ctx, departingID)
|
|
Expect(err).To(MatchError(clustersvc.ErrInstanceNotFound))
|
|
})
|
|
|
|
It("reports a killed replica as unreachable, reaps what it owned, and evicts no worker", func() {
|
|
// This is the absence rule, pinned before phase 2 can depend on it. A
|
|
// wrong implementation lets a peer that will not answer surface as node
|
|
// absence, and a caller entitled to act on absence then reclaims what
|
|
// the peer was running: a network hiccup between two healthy replicas
|
|
// evicts healthy workers.
|
|
//
|
|
// It also pins the reaper: the connection rows a dead replica owned are
|
|
// swept by the same sweeper that decides the replica is dead, so the
|
|
// two can never disagree about who is alive.
|
|
c, dsn := startClusterOnFreshDB(2, 1)
|
|
|
|
client, err := c.AdminSession(0)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
|
|
// The worker registers with frontend 0, so frontend 1 is the replica
|
|
// that can die without taking the worker's registrar with it.
|
|
registrar, err := c.WorkerRegistrar(0)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
Expect(registrar).To(Equal(0), "this spec kills frontend 1 and needs the worker to have registered elsewhere")
|
|
|
|
probe := newRosterProbe(c, client, 0)
|
|
Eventually(probe.healthyNames, nodeRosterTimeout, nodeRosterPoll).
|
|
Should(ContainElement(c.WorkerName(0)), probe.describe)
|
|
workerID := probe.idOf(c.WorkerName(0))
|
|
Expect(workerID).ToNot(BeEmpty())
|
|
|
|
roster := newInstanceRoster(openClusterDB(dsn))
|
|
awaitReplicas(roster, hostPortOf(c.FrontendURL(0)), hostPortOf(c.FrontendURL(1)))
|
|
survivorID := roster.idAt(hostPortOf(c.FrontendURL(0)))
|
|
doomedID := roster.idAt(hostPortOf(c.FrontendURL(1)))
|
|
Expect(survivorID).ToNot(BeEmpty())
|
|
Expect(doomedID).ToNot(BeEmpty())
|
|
|
|
// Give frontend 1 the worker's tunnel. Phase 2 makes the worker do this
|
|
// by dialling; here the claim is written directly, because the point
|
|
// under test is what happens to the claim when its owner dies.
|
|
ctx := context.Background()
|
|
epoch, err := roster.registry.Claim(ctx, workerID, doomedID)
|
|
Expect(err).ToNot(HaveOccurred())
|
|
Expect(epoch).ToNot(BeZero())
|
|
|
|
Expect(c.KillFrontend(1)).To(Succeed())
|
|
Eventually(func() bool { return c.FrontendAlive(1) }, "20s", "500ms").Should(BeFalse())
|
|
|
|
// The row is still there for the whole liveness window, so this is the
|
|
// case that matters: the peer is KNOWN and will not answer.
|
|
dialCtx, cancel := context.WithTimeout(ctx, peerDialTimeout)
|
|
defer cancel()
|
|
pool := clustersvc.NewPeerPool("e2e-peer", c.RegistrationToken(), joinAsPeer(roster, "e2e-peer"), roster.registry)
|
|
DeferCleanup(pool.Close)
|
|
_, err = pool.Open(dialCtx, doomedID)
|
|
Expect(err).To(MatchError(clustersvc.ErrPeerUnreachable))
|
|
Expect(err).ToNot(MatchError(clustersvc.ErrInstanceNotFound),
|
|
"a dead replica whose row is still present is unreachable, not absent")
|
|
|
|
// The survivor sweeps the dead replica and, in the same pass, the claim
|
|
// it left behind.
|
|
Eventually(roster.addresses, deadReplicaTimeout, instanceRosterPoll).
|
|
Should(ConsistOf(hostPortOf(c.FrontendURL(0))), roster.describe)
|
|
ownerErr := func() error {
|
|
_, _, err := roster.registry.OwnerRow(ctx, workerID)
|
|
return err
|
|
}
|
|
Eventually(ownerErr, deadReplicaTimeout, instanceRosterPoll).
|
|
Should(MatchError(clustersvc.ErrNoConnection),
|
|
"the claim held by a replica that no longer exists was never reaped")
|
|
|
|
// And the worker survives the sweep that removed its owner. This is a
|
|
// window after the reaping, not a watch over the whole scenario:
|
|
// Consistently starts here, so what it rules out is the sweep, or
|
|
// anything reacting to it, taking the worker with it.
|
|
Consistently(probe.healthyNames, "6s", "1s").
|
|
Should(ContainElement(c.WorkerName(0)), probe.describe)
|
|
})
|
|
})
|