The controller's machine moves it from the container to a process by starting the process first and removing the container once the process is up (mesh-host's `replaces`). For that moment two controllers share the store and the bus. Checked what each does: - the seat's verbs: a queue group per seat, each call answered once. Safe. - the controller's consumers on CONTROL and EVENTS: push consumers with no delivery group, so the second bind is refused with "consumer is already bound" and serve exited. The process would restart for ever, the host would never see it up, and the container would never go. The second controller now stands by and binds when the first lets go (tested on a real bus; fails without the change). - plans: read, changed and saved whole by the 30s timer, by build outcomes, by a merge and by `plans stop`. Two timers would each ask a tier the other had just asked. Working the plans now takes a session-level advisory lock on the inventory: the timer skips while another holds it, the other paths wait for it. Build asks happen only inside plan work and are covered by the same lock.
149 lines
5.2 KiB
Go
149 lines
5.2 KiB
Go
package inventory
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"slices"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// ErrNodeBusy is a node another act kept holding for longer than a caller waits.
|
|
var ErrNodeBusy = errors.New("another act is composing or sending that node's declaration")
|
|
|
|
// HoldWaitFor is how long HoldNodes waits for a node another act holds, and HoldPoll how often it
|
|
// looks again. Variables so a test need not wait minutes.
|
|
var (
|
|
HoldWaitFor = 2 * time.Minute
|
|
HoldPoll = 250 * time.Millisecond
|
|
)
|
|
|
|
// HoldNodes serialises composing and sending a declaration per node (novox/hq ADR 0100): while
|
|
// one caller holds a node, another asking for it waits. Without it a push that composed a node as
|
|
// adopted could send that declaration after `converge --yes` sent the converged one, and the node
|
|
// would return to adopted with nobody having asked.
|
|
//
|
|
// Session-level advisory locks on one connection, all or none: a set not wholly free is given back
|
|
// at once, so two callers holding overlapping sets never each wait on the other. **A waiter pins
|
|
// no connection.** It looks again every HoldPoll with a connection borrowed for the look, and gives
|
|
// up after HoldWaitFor with ErrNodeBusy naming the node — so a stuck holder costs the pool one
|
|
// connection, never one per caller queued behind it. Release gives every one back, and may be
|
|
// called more than once.
|
|
func (i *Inventory) HoldNodes(ctx context.Context, names []string) (func(), error) {
|
|
sorted := slices.Clone(names)
|
|
slices.Sort(sorted)
|
|
sorted = slices.Compact(sorted)
|
|
deadline := time.Now().Add(HoldWaitFor)
|
|
for {
|
|
release, busy, err := i.tryHold(ctx, sorted)
|
|
if err != nil || busy == "" {
|
|
return release, err
|
|
}
|
|
if time.Now().After(deadline) {
|
|
return nil, fmt.Errorf("%w: %s has been held for over %s — try again once it is done",
|
|
ErrNodeBusy, busy, HoldWaitFor)
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err()
|
|
case <-time.After(HoldPoll):
|
|
}
|
|
}
|
|
}
|
|
|
|
// tryHold takes every named node's lock or none, and says which node was busy when it took none.
|
|
func (i *Inventory) tryHold(ctx context.Context, sorted []string) (func(), string, error) {
|
|
conn, err := i.store.Pool().Acquire(ctx)
|
|
if err != nil {
|
|
return nil, "", err
|
|
}
|
|
var once sync.Once
|
|
release := func() {
|
|
once.Do(func() {
|
|
// Unlocking all of this session's advisory locks, then handing the connection back: a
|
|
// connection returned still holding one would hold it for whoever borrows it next.
|
|
_, err := conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock_all()`)
|
|
if err != nil {
|
|
// The session's locks die with the session: close it rather than return it.
|
|
_ = conn.Conn().Close(context.WithoutCancel(ctx))
|
|
}
|
|
conn.Release()
|
|
})
|
|
}
|
|
for _, name := range sorted {
|
|
var took bool
|
|
if err := conn.QueryRow(ctx,
|
|
`select pg_try_advisory_lock(hashtext('mesh-node-declaration:' || $1)::bigint)`,
|
|
name).Scan(&took); err != nil {
|
|
release()
|
|
return nil, "", err
|
|
}
|
|
if !took {
|
|
release()
|
|
return nil, strings.TrimSpace(name), nil
|
|
}
|
|
}
|
|
return release, "", nil
|
|
}
|
|
|
|
// ErrPlansBusy is the plans held by another act — on a machine replacing its controller, the other
|
|
// controller — for longer than a caller waits, or at all for one that does not wait.
|
|
var ErrPlansBusy = errors.New("another controller is working the plans")
|
|
|
|
// HoldPlans makes working the plans one act at a time, across every controller on the store
|
|
// (novox/hq issue 213). A plan is read, changed and written whole; two controllers doing that at
|
|
// once — the old and the new for the moment a machine hands its controller over, or a controller
|
|
// and a person's `plans stop` — each act on what the other has not saved yet: a tier asked twice,
|
|
// an outcome written over. A session-level advisory lock on one connection, released by the
|
|
// returned function and by the session ending, so a controller that dies holding it holds nothing.
|
|
//
|
|
// wait false gives ErrPlansBusy at once when another holds them — the timer's way: the holder is
|
|
// moving the plans already. wait true looks again every HoldPoll for up to HoldWaitFor — an
|
|
// outcome's or a merge's way, which must be written.
|
|
func (i *Inventory) HoldPlans(ctx context.Context, wait bool) (func(), error) {
|
|
deadline := time.Now().Add(HoldWaitFor)
|
|
for {
|
|
release, took, err := i.tryLock(ctx, "mesh-plans")
|
|
if err != nil || took {
|
|
return release, err
|
|
}
|
|
if !wait || time.Now().After(deadline) {
|
|
return nil, ErrPlansBusy
|
|
}
|
|
select {
|
|
case <-ctx.Done():
|
|
return nil, ctx.Err()
|
|
case <-time.After(HoldPoll):
|
|
}
|
|
}
|
|
}
|
|
|
|
// tryLock takes one named advisory lock on a connection of its own, or gives the connection back.
|
|
func (i *Inventory) tryLock(ctx context.Context, key string) (func(), bool, error) {
|
|
conn, err := i.store.Pool().Acquire(ctx)
|
|
if err != nil {
|
|
return nil, false, err
|
|
}
|
|
var once sync.Once
|
|
release := func() {
|
|
once.Do(func() {
|
|
if _, err := conn.Exec(context.WithoutCancel(ctx), `select pg_advisory_unlock_all()`); err != nil {
|
|
_ = conn.Conn().Close(context.WithoutCancel(ctx))
|
|
}
|
|
conn.Release()
|
|
})
|
|
}
|
|
var took bool
|
|
if err := conn.QueryRow(ctx, `select pg_try_advisory_lock(hashtext($1)::bigint)`, key).Scan(&took); err != nil {
|
|
release()
|
|
return nil, false, err
|
|
}
|
|
if !took {
|
|
release()
|
|
return nil, false, nil
|
|
}
|
|
return release, true, nil
|
|
}
|