Received: from malur.postgresql.org ([217.196.149.56]) by arkaria.postgresql.org with esmtps (TLS1.3:ECDHE_RSA_AES_256_GCM_SHA384:256) (Exim 4.92) (envelope-from ) id 1jLbX4-0001Ps-Va for pgsql-hackers@arkaria.postgresql.org; Mon, 06 Apr 2020 23:51:59 +0000 Received: from localhost ([127.0.0.1] helo=malur.postgresql.org) by malur.postgresql.org with esmtp (Exim 4.92) (envelope-from ) id 1jLbX3-0006gH-OW for pgsql-hackers@arkaria.postgresql.org; Mon, 06 Apr 2020 23:51:57 +0000 Received: from magus.postgresql.org ([2a02:c0:301:0:ffff::29]) by malur.postgresql.org with esmtps (TLS1.3:ECDHE_RSA_AES_256_GCM_SHA384:256) (Exim 4.92) (envelope-from ) id 1jLbX3-0006gA-GJ for pgsql-hackers@lists.postgresql.org; Mon, 06 Apr 2020 23:51:57 +0000 Received: from mail-qt1-x842.google.com ([2607:f8b0:4864:20::842]) by magus.postgresql.org with esmtps (TLS1.3:ECDHE_RSA_AES_128_GCM_SHA256:128) (Exim 4.92) (envelope-from ) id 1jLbWz-0006pA-M4 for pgsql-hackers@lists.postgresql.org; Mon, 06 Apr 2020 23:51:56 +0000 Received: by mail-qt1-x842.google.com with SMTP id n17so1380992qtv.4 for ; Mon, 06 Apr 2020 16:51:53 -0700 (PDT) DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=2ndquadrant-com.20150623.gappssmtp.com; s=20150623; h=date:from:to:cc:subject:message-id:mime-version:content-disposition :content-transfer-encoding:in-reply-to:user-agent; bh=vLrx/k7CHDpPVMznliNj4PjCiNs6CsFHufFwKzwxDCY=; b=taKSNzU1sYdT32jC0WpG00qqwdAbS3gGPKJ5K8BvpTP64r+iWc6ufsxijF5mHyZFot /FpLPt6Q7S2B/8+2/hLOl+GHgEblMpF8xBlJ3SNDqntPBOMDjsThXNT7pE+b0NENIJzM VcD+hu32oFLx6XrApisUWQFHdYPuCOQIYUf8yBnX/5IbnR9jOBHLlrzZl5NG6T9HGS+8 AjtOhu4adQ7T50VuCOmlZyqXi6B+8goXukdwHFHwc5xIzFNMW9/qHKbRFJAsWNW54dZG yj/id3myD4Q/Km/jzWU8XvzN0K6xuftrwrVguavldGuxNLpOqF7JRvAm3HqlEtWkmKIv TwgQ== X-Google-DKIM-Signature: v=1; a=rsa-sha256; c=relaxed/relaxed; d=1e100.net; s=20161025; h=x-gm-message-state:date:from:to:cc:subject:message-id:mime-version :content-disposition:content-transfer-encoding:in-reply-to :user-agent; bh=vLrx/k7CHDpPVMznliNj4PjCiNs6CsFHufFwKzwxDCY=; b=mKAMljDvfelkRNAxThXl9imN6qeJjLDfVRSNbYo8MUPv0oVvK87meQI83dPBxPDfdg BlZRr+Dbab9WOwYw0eTTF8fnrAQrWiGOKDjT/eBfQbYlGOM7UI1AEZi6MiiX0ogzgAAD sRCGyW+KDXcOZd1B8zP63bl01BIzpSG6Q+9ZutZkXlWbVkKY3YzEpWcrsVIX+Qsb2RmE x848O5K81OHhHm/11VLI6HIBNR6xZmVFu/apFT9PeXHKmXr0YBK0xJ3+/Ky7YjnYClvR wcSayCsPjGSyC5kw30GF2U+a4B7iUz477y2VRKFBSiEOxssPmoxgXu9ZE3WS4g6Y8b3/ DNzg== X-Gm-Message-State: AGi0PuYBNevheQfU1uh9eqxxwUp11jAQxbzwEZ8p6tAWhkgM7myTEsDB 2mQ2s+iJH9VG/rg6Xl/CfJvyMA== X-Google-Smtp-Source: APiQypLx5lVtpej2G824isP7gapOQOIi+DQLfFRFjFH5TYjkzbRusbchiQm13yq5PA/nwqUZlKm04A== X-Received: by 2002:ac8:1757:: with SMTP id u23mr2138365qtk.138.1586217111499; Mon, 06 Apr 2020 16:51:51 -0700 (PDT) Received: from nimloth.alvh.no-ip.org ([190.95.18.252]) by smtp.gmail.com with ESMTPSA id t15sm1415493qtc.64.2020.04.06.16.51.50 (version=TLS1_3 cipher=TLS_AES_256_GCM_SHA384 bits=256/256); Mon, 06 Apr 2020 16:51:50 -0700 (PDT) Received: by nimloth.alvh.no-ip.org (Postfix, from userid 1000) id 6A50E300A04; Mon, 6 Apr 2020 19:51:48 -0400 (-04) Date: Mon, 6 Apr 2020 19:51:48 -0400 From: Alvaro Herrera To: Kyotaro Horiguchi Cc: jgdr@dalibo.com, andres@anarazel.de, michael@paquier.xyz, sawada.mshk@gmail.com, peter.eisentraut@2ndquadrant.com, pgsql-hackers@lists.postgresql.org, thomas.munro@enterprisedb.com, sk@zsrv.org, michael.paquier@gmail.com Subject: Re: [HACKERS] Restricting maximum keep segments by repslots Message-ID: <20200406235148.GA29998@alvherre.pgsql> MIME-Version: 1.0 Content-Type: multipart/mixed; boundary="DocE+STaALJfprDB" Content-Disposition: inline Content-Transfer-Encoding: 8bit In-Reply-To: <20200406221555.GA16211@alvherre.pgsql> User-Agent: Mutt/1.10.1 (2018-07-13) List-Id: List-Help: List-Subscribe: List-Post: List-Owner: List-Archive: Precedence: bulk --DocE+STaALJfprDB Content-Type: text/plain; charset=iso-8859-1 Content-Disposition: inline Content-Transfer-Encoding: 8bit On 2020-Apr-06, Alvaro Herrera wrote: > I think there's a race condition in this: if we kill a walsender and it > restarts immediately before we (checkpoint) can acquire the slot, we > will wait for it to terminate on its own. Fixing this requires changing > the ReplicationSlotAcquire API so that it knows not to wait but not > raise error either (so we can use an infinite loop: "acquire, if busy > send signal") I think this should do it, but I didn't test it super-carefully and the usage of the condition variable is not entirely kosher. -- Álvaro Herrera https://www.2ndQuadrant.com/ PostgreSQL Development, 24x7 Support, Remote DBA, Training & Services --DocE+STaALJfprDB Content-Type: text/x-diff; charset=us-ascii Content-Disposition: attachment; filename="v25delta-walsender-kill-loop.patch" commit af47f2e05316e7e89ee7e5b59b5f5fe6ae421508 Author: Alvaro Herrera AuthorDate: Mon Apr 6 19:29:56 2020 -0400 CommitDate: Mon Apr 6 19:51:06 2020 -0400 loop in InvalidateObsoleteReplicationSlots fixes a race condition on slot acquisition diff --git a/src/backend/replication/logical/logicalfuncs.c b/src/backend/replication/logical/logicalfuncs.c index 04510094a8..f5384f1df8 100644 --- a/src/backend/replication/logical/logicalfuncs.c +++ b/src/backend/replication/logical/logicalfuncs.c @@ -225,7 +225,7 @@ pg_logical_slot_get_changes_guts(FunctionCallInfo fcinfo, bool confirm, bool bin else end_of_wal = GetXLogReplayRecPtr(&ThisTimeLineID); - ReplicationSlotAcquire(NameStr(*name), true); + (void) ReplicationSlotAcquire(NameStr(*name), SAB_Error); PG_TRY(); { diff --git a/src/backend/replication/slot.c b/src/backend/replication/slot.c index 31e12e4043..e7960cb48e 100644 --- a/src/backend/replication/slot.c +++ b/src/backend/replication/slot.c @@ -325,9 +325,15 @@ ReplicationSlotCreate(const char *name, bool db_specific, /* * Find a previously created slot and mark it as used by this backend. + * + * The return value is only useful if behavior is SAB_Inquire, in which + * it's zero if we successfully acquired the slot, or the PID of the + * owning process otherwise. If behavior is SAB_Error, then trying to + * acquire an owned slot is an error. If SAB_Block, we sleep until the + * slot is released by the owning process. */ -void -ReplicationSlotAcquire(const char *name, bool nowait) +int +ReplicationSlotAcquire(const char *name, SlotAcquireBehavior behavior) { ReplicationSlot *slot; int active_pid; @@ -392,11 +398,13 @@ retry: */ if (active_pid != MyProcPid) { - if (nowait) + if (behavior == SAB_Error) ereport(ERROR, (errcode(ERRCODE_OBJECT_IN_USE), errmsg("replication slot \"%s\" is active for PID %d", name, active_pid))); + else if (behavior == SAB_Inquire) + return active_pid; /* Wait here until we get signaled, and then restart */ ConditionVariableSleep(&slot->active_cv, @@ -412,6 +420,9 @@ retry: /* We made this slot active, so it's ours now. */ MyReplicationSlot = slot; + + /* success */ + return 0; } /* @@ -518,7 +529,7 @@ ReplicationSlotDrop(const char *name, bool nowait) { Assert(MyReplicationSlot == NULL); - ReplicationSlotAcquire(name, nowait); + (void) ReplicationSlotAcquire(name, nowait ? SAB_Error : SAB_Block); ReplicationSlotDropAcquired(); } @@ -1097,7 +1108,6 @@ restart: ReplicationSlot *s = &ReplicationSlotCtl->replication_slots[i]; XLogRecPtr restart_lsn = InvalidXLogRecPtr; char *slotname; - int wspid; if (!s->in_use) continue; @@ -1112,21 +1122,27 @@ restart: slotname = pstrdup(NameStr(s->data.name)); restart_lsn = s->data.restart_lsn; - wspid = s->active_pid; SpinLockRelease(&s->mutex); LWLockRelease(ReplicationSlotControlLock); - if (wspid != 0) + for (;;) { - ereport(LOG, - (errmsg("terminating walsender %d because replication slot is too far behind", - wspid))); - (void) kill(wspid, SIGTERM); - } + int wspid = ReplicationSlotAcquire(slotname, SAB_Inquire); - /* Here we wait until the walsender is gone */ - ReplicationSlotAcquire(slotname, false); + /* no walsender? success! */ + if (wspid == 0) + break; + + ereport(LOG, + (errmsg("terminating walsender %d because replication slot \"%s\" is too far behind", + wspid, slotname))); + (void) kill(wspid, SIGTERM); + + ConditionVariableTimedSleep(&s->active_cv, + 10, WAIT_EVENT_REPLICATION_SLOT_DROP); + ConditionVariableCancelSleep(); + } ereport(LOG, (errmsg("invalidating slot \"%s\" because its restart_lsn %X/%X exceeds max_slot_wal_keep_size", diff --git a/src/backend/replication/slotfuncs.c b/src/backend/replication/slotfuncs.c index 91a5d0f290..f8336129d9 100644 --- a/src/backend/replication/slotfuncs.c +++ b/src/backend/replication/slotfuncs.c @@ -592,7 +592,7 @@ pg_replication_slot_advance(PG_FUNCTION_ARGS) moveto = Min(moveto, GetXLogReplayRecPtr(&ThisTimeLineID)); /* Acquire the slot so we "own" it */ - ReplicationSlotAcquire(NameStr(*slotname), true); + (void) ReplicationSlotAcquire(NameStr(*slotname), SAB_Error); /* A slot whose restart_lsn has never been reserved cannot be advanced */ if (XLogRecPtrIsInvalid(MyReplicationSlot->data.restart_lsn)) diff --git a/src/backend/replication/walsender.c b/src/backend/replication/walsender.c index 9e5611574c..06e8b79036 100644 --- a/src/backend/replication/walsender.c +++ b/src/backend/replication/walsender.c @@ -595,7 +595,7 @@ StartReplication(StartReplicationCmd *cmd) if (cmd->slotname) { - ReplicationSlotAcquire(cmd->slotname, true); + (void) ReplicationSlotAcquire(cmd->slotname, SAB_Error); if (SlotIsLogical(MyReplicationSlot)) ereport(ERROR, (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), @@ -1132,7 +1132,7 @@ StartLogicalReplication(StartReplicationCmd *cmd) Assert(!MyReplicationSlot); - ReplicationSlotAcquire(cmd->slotname, true); + (void) ReplicationSlotAcquire(cmd->slotname, SAB_Error); /* * Force a disconnect, so that the decoding code doesn't need to care diff --git a/src/include/replication/slot.h b/src/include/replication/slot.h index 6e469ea749..f984bfd7a6 100644 --- a/src/include/replication/slot.h +++ b/src/include/replication/slot.h @@ -36,6 +36,14 @@ typedef enum ReplicationSlotPersistency RS_TEMPORARY } ReplicationSlotPersistency; +/* For ReplicationSlotAcquire, q.v. */ +typedef enum SlotAcquireBehavior +{ + SAB_Error, + SAB_Block, + SAB_Inquire +} SlotAcquireBehavior; + /* * On-Disk data of a replication slot, preserved across restarts. */ @@ -184,7 +192,7 @@ extern void ReplicationSlotCreate(const char *name, bool db_specific, extern void ReplicationSlotPersist(void); extern void ReplicationSlotDrop(const char *name, bool nowait); -extern void ReplicationSlotAcquire(const char *name, bool nowait); +extern int ReplicationSlotAcquire(const char *name, SlotAcquireBehavior behavior); extern void ReplicationSlotRelease(void); extern void ReplicationSlotCleanup(void); extern void ReplicationSlotSave(void); --DocE+STaALJfprDB--