flyteorg / flyteorg/flyte

perf: don't block action writes on the NOTIFY pump

Open
#7,757 2 comments 0 reactions 1 assignee Claimed by @0yukali0 View on GitHub
enhancement flyte2 help wanted scale
Dominant language
Go
Stars
7.5k
Forks
886
Avg merge
1d 12h
Merged PRs (30d)
120

Description

## Why

A PostgreSQL `NOTIFY` in the runs-service is a **wakeup hint** — the payload carries only an action
identifier, and every listener re-reads state from the database when it wakes. Losing one is
recoverable. But today, sending one can block a write RPC indefinitely.

[`notifyActionUpdate`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1111-L1120) does a **blocking** channel send:

```go
select {
case r.actionNotifyCh <- payload:
case <-ctx.Done():
logger.Errorf(ctx, "Action NOTIFY send cancelled for %s: %v", payload, ctx.Err())
}
```

The channel is [buffered at **256**](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L54-L55) and drained by a single goroutine that round-trips to
Postgres once per payload — a ceiling of roughly **240 notifications/second** (measured: 4.1 ms per
`pg_notify` round trip from the same VPC, versus 32 µs for the `pg_notify` call itself). When demand
exceeds that, the buffer fills and the send blocks — **inside the RPC handler**, because
[`UpdateActionPhase` calls it](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L482-L485) right after the row is written.

### What is *not* at risk: the database write

Worth stating up front, because it is the first question everyone asks. The status update is
**already committed before the notification is attempted** — every one of the seven notify call
sites runs after its DB statement returns, and [`UpdateActionPhase`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L429) additionally
gates on `rowsAffected > 0`:

```go
result, err := r.db.ExecContext(ctx, queryBuilder.String(), args...) // committed here
if err != nil { return err } // no notify on failure
rowsAffected, err := result.RowsAffected()
if rowsAffected > 0 {
r.notifyActionUpdate(ctx, actionID) // only after the row is written
}
```

These are autocommit single statements — there is no transaction spanning the write and the
notification. So blocking on the channel does **not** make the write more likely to succeed; it
only delays the RPC's return after the write already succeeded.

What *is* at risk is watcher liveness: `WatchActionUpdates` / `WatchAllActionUpdates` take their
wakeup from this notification and re-read state from the database when woken. Lose a wakeup and a
stream — and therefore the phase a user sees in the UI — can lag until something else touches that
action. So the fix must not block the write path **and** must not lose a wakeup. The design below
does both; see "What to change".

The result is that watch fan-out latency becomes write latency. Measured on a dev cluster under a
100,000-action benchmark:

| RPC | rate | p50 | p95 | mean |
|---|---|---|---|---|
| `InternalRunService.UpdateActionStatus` | **101.8/s** | 394 ms | **3085 ms** | 780 ms |
| `InternalRunService.RecordActionEvents` | 25.6/s | 652 ms | 2418 ms | 800 ms |
| `ActionsService.Enqueue` | 57.5/s | — | 1351 ms | 375 ms |

Meanwhile the database was idle (`WriteLatency` 0.9 ms, **14 connections** open out of 100, all in
`ClientRead`) and the pod used 1.34 of 4 cores with 0.09% CPU throttling. About **78 requests were
in flight** (Little's law: 101.8 rps × 0.78 s) waiting on neither CPU nor Postgres.

Batching the pump's round trips (see the companion issue) raises the ceiling, but the structural
problem remains: **a durable write should not wait on a best-effort notification.** As long as the
send can block, any future slowdown in the notify path silently becomes user-visible write latency.

## What to change

The fix must satisfy two constraints at once: **never block the write path**, and **never lose a
wakeup** — a lost wakeup means a run sits at a stale phase in the UI until something else happens
to that action.

Both are achievable, because the notification payload is an *identity*
(`project/domain/run/action`), not a state. Listeners re-read from the database when woken. So
collapsing five notifications for the same action into one is **not** a partial delivery — the
single wakeup still causes a read of the latest state. What must never be lost is the *last*
wakeup for a given action.

That points at replacing the FIFO channel with a **coalescing pending set**:

```go
// today:
actionNotifyCh chan string // 256 FIFO slots, blocking send

// proposed:
pendingActions map[string]struct{} // mutex-guarded; key = the payload
pendingCh chan struct{} // capacity 1, non-blocking "something is pending" signal
```

(Name the signal channel to match the surrounding style — `pendingCh` alongside `pendingActions`
reads as "what is pending" plus "a nudge that something is pending". The existing fields use the
`...Ch` suffix, so `wakeCh`/`signalCh` would also be fine; just avoid a bare `wake`.)

- **Writer** ([`notifyActionUpdate`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1111-L1120)): take the lock, insert the key, unlock, then a
non-blocking send on `pendingCh`. It cannot block and cannot discard a key. A signal that is
already buffered needs no second one — the map is the queue, the channel is only a nudge.
- **Pump** ([`runNotifyLoop`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1039)): on each signal, swap the map out under the lock, then emit
the whole set — ideally as one batched statement, which is exactly what the companion batching
issue does. The two changes compose.

### Delivery: keep the key until the NOTIFY succeeds

"If the writer doesn't block, how do we know the NOTIFY was actually sent?" — worth answering
directly, because **blocking does not answer it today either**. [`execNotify`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1062-L1076)
discards the payload on any failure:

```go
if conn == nil {
logger.Errorf(ctx, "No NOTIFY connection available, dropping %s notification", channel)
return // payload lost
}
if _, err := conn.ExecContext(ctx, "SELECT pg_notify($1, $2)", channel, payload); err != nil {
logger.Errorf(ctx, "Failed to NOTIFY %s: %v", channel, err)
if isConnError(err) { reconnect() }
// payload lost — no retry
}
```

The blocking send only guarantees that the payload reached a 256-slot buffer. Once the pump takes
it, a connection blip drops it silently. So delivery is best-effort today, and a run can already
sit at a stale phase in the UI after a transient database hiccup.

The pending set fixes this, and that is the main reason to prefer it over a bigger buffer:

1. The pump swaps the pending set out into a batch under the lock.
2. It attempts the batched `pg_notify`.
3. On success the keys are gone.
4. **On failure it merges the keys back into the pending set**, reconnects, and retries with backoff.

Step 4 is safe *because it is a set*: re-adding a key is idempotent, so a retry cannot duplicate or
double-count, and any new updates for the same action that arrived meanwhile merge into the same
entry. A FIFO would force a choice between dropping and unbounded re-queueing. The result is
**at-least-once delivery for the life of the process** — the right target, since a notification is
an idempotent wakeup and a duplicate costs one redundant re-read.

**What about a runs-service restart?** In-memory pending entries are lost, but that does not strand
the UI, because every new watch stream starts from the database rather than from notifications.
[`WatchActions`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/service/run_service.go#L1336) subscribes first, *then* sends a full snapshot paged out of the
`actions` table, and only then streams deltas:

```go
// Start watching for updates from DB first to prevent event miss
go s.repo.ActionRepo().WatchAllActionUpdates(ctx, runID, updatesCh, errsCh)
...
if err := s.listAndSendAllActions(ctx, runID, rsm, stream); err != nil { // full snapshot
```

So a restart breaks the streams, clients reconnect, and each reconnect re-reads authoritative state
that was committed long before any notification was attempted. Losing in-flight wakeups costs a
reconnect window, not a permanently stale run. This is also *why* the notification path can be
best-effort at all: it is an optimization over polling, and Postgres is the source of truth.

Why this is better than a bigger buffer or a drop policy:

| | today | drop when full | pending set |
|---|---|---|---|
| Blocks the write path | **yes** | no | no |
| Can lose a wakeup | yes, on `ctx.Done()` | yes, under pressure | **no** |
| Queue bounded by | 256 updates | 256 updates | distinct *actions* |

The memory bound is the key property: the set can never hold more entries than there are distinct
actions with an unsent notification, no matter how fast updates arrive. At 65,000 live actions
that is a few MB.

It also removes a loss path that exists **today**: the current `select` discards the notification
on `ctx.Done()`, so a client disconnecting mid-request can lose a wakeup that other watchers
needed. A non-blocking send has no reason to consult the request context at all.

Files:

- `runs/repository/impl/action.go` — [`notifyActionUpdate`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1111-L1120),
[`notifyRunUpdate`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L952-L961), the [channel declarations](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L40-L41) and their
[construction](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L54-L55), [`runNotifyLoop`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1039) and [`drainAndExec`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1078-L1090).
- `runs/repository/impl/action_test.go` — cover the properties below.

Ordering note: `notifyRunUpdate` needs the same treatment, and the two sets can share one pump.

## Outcome

- [ ] `UpdateActionPhase` never blocks on notification delivery, however slow or stalled the pump is
- [ ] No wakeup is lost: for every action that was updated, at least one notification is delivered
after its last update
- [ ] A failed `pg_notify` retries rather than dropping the payload — keys stay pending until the
statement succeeds (strictly better than today's behaviour)
- [ ] Repeated notifications for the same action between drains collapse into one
- [ ] Request-context cancellation no longer discards a notification
- [ ] Memory is bounded by distinct pending actions rather than by update volume
- [ ] Tests cover: pump stalled → writer returns immediately and the row is still written;
N updates to one action between drains → one notification; every updated action is
represented in the delivered set

## Getting started

- **Reproduce:** generate ~100 action-updates/second against a dev cluster and watch
`rpc_server_duration_milliseconds` for `UpdateActionStatus`. The
[flyteorg/benchmark](https://github.com/flyteorg/benchmark) suite does this with
`uv run scripts/v2/swarm.py --k 25 --n 2000`.
- **Relevant code:** [`notifyActionUpdate`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1111-L1120), then the pump in
[`runNotifyLoop`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L1039), then the listener side in [`processNotifications`](https://github.com/flyteorg/flyte/blob/a2d7e82be425372e6d94252932b31800b34992f3/runs/repository/impl/action.go#L901).
- **How to test:** `go test ./runs/repository/impl/...` — needs a local Postgres on port **15432**
(`user/password: postgres`, db `flyte_runs_test`; see `testDbConfig` in `action_test.go`).
- **Setup:** [CONTRIBUTING.md](https://github.com/flyteorg/flyte/blob/main/CONTRIBUTING.md)

Related — three issues on the same bottleneck:
- #7756 batch the pump's round trips
- #7757 don't block action writes on the pump (this one)
- #7758 instrument the pump

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.