github.com-googleapis-google-cloud-go
all · 27 devs · built 2026-08-09
Repository snapshot
Monthly reports
No monthly reports available yet.
Performance over time
ETV stacked by Growth, Maintenance and Fixes — 90-day moving average, normalized to ETV / month.
Average performance per developer
ETV per active developer per month — 30-day moving average.
Active developers over time
Unique developers committing each day — 90-day moving average.
Knowledge concentration
How dependent is this repo on a small number of contributors? Higher top-1 share = higher key-person risk.
Baha Aiman owns 15.1 % of commits.
Top contributors
Most impactful commits
Top 20 by ETV in the all-time window.
- 6.0ETVfeat(spanner): complete location-aware routing resilience and observability (#14418) ## Summary This change expands the location-aware routing path in the Spanner Go client with endpoint lifecycle management, stronger endpoint health handling, retry-aware exclusion, latency-aware replica selection, and route-selection tracing. ## What Changed - enabled location-aware routing automatically for experimental-host configurations while preserving env-var override behavior - added endpoint cache support for GetIfPresent, Evict, and DefaultChannel - switched endpoint health evaluation to real gRPC connectivity state, including transient-failure detection - added an endpoint lifecycle manager with: - background probing - idle eviction - endpoint recreation requests - recent transient-failure eviction tracking - integrated lifecycle handling into routing and affinity fallback paths - added request-id keyed one-shot endpoint exclusion for RESOURCE_EXHAUSTED retries - expanded skipped_tablet_uid reporting for transient failures and recent transient-failure evictions, with deduplication - added OpenTelemetry route-selection span attributes and events - updated focused tests and goldens for the new routing behavior ## Behavioral Notes - stopped or transient-failure endpoints are now treated differently from simply non-ready endpoints during tablet selection - retry attempts can avoid previously selected routed endpoints after RESOURCE_EXHAUSTED - leader-preferred requests still favor close leaders, but can fall back to a closer non-leader replica when the leader is too far away - route selection is now observable through tracing attributes and eventsrahul2393 · 77aa4df8 · 2026-04-14
- 5.3ETVfeat(workloadmanager): add new clients (#13849) PiperOrigin-RevId: 869327994shollyman · 13288723 · 2026-02-14
- 4.6ETVfeat(bigtable): add SessionPoolImpl (two-tier pool + scaling + debug) (#20225) ## Summary Second of five PRs porting the session pool infrastructure from `feat/bigtable-sessionz-debug` (which powers the recycle-repro fleet). Adds `SessionPoolImpl`: the concrete two-tier read/write session pool for one resource. ~4000 LOC across five files plus matching tests. ## Stack - [ ] PR-1: sessionList (#20224) — per-AFE bucketing data structure. **Not yet merged.** - [x] **PR-2 (this)** — SessionPoolImpl (pool + scaling + debug + snapshot). - [ ] PR-3 — `SessionPool` / `Invoker` interfaces + `sessionClient` / `sessionTable` factory wiring. - [ ] PR-4 — Debug pages (`sessionz` / `afez` / `flightz` / `loadz` under `bigtable/debugview/`). - [ ] PR-5 — `bigtable.Client` integration + release notes. **Because PR-1 has not landed, this PR is opened against `main` and the diff includes PR-1's commits.** Once #20224 merges the base can be re-targeted (or this branch rebased) so only PR-2's own delta shows. ## What lands **SessionPoolImpl** (5 files, ~2400 LOC prod + ~2200 LOC tests): - `session_pool.go` — struct + constructor + `Invoke` + `CheckoutSession` (waiter queue, deadline propagation) + pluggable picker via the AFE picker from #20204. - `session_pool_lifecycle.go` — `SessionHooks` wiring, consecutive-failure breaker, `Close` (5-phase teardown), `WaitGoroutines` / `spawns.Wait` choreography so no session-owned goroutine outlives the pool. - `session_pool_scaling.go` — `Tick` loop, `createSession` (dial + `OpenSession` + hook registration), `pendingStarts` / `startingSessions` accounting so scale-up decisions never double-count in-flight opens. Uses the channel-pool pick hint (`ChannelPickHintInto`, added to `connpool.go`) to attribute each session to its underlying channel. - `session_pool_debug.go` — `PoolSnapshot` / slow-vRPC ring / per-close-reason counters / scaling-history buffer / `pickHistory` ring — the input to the sessionz / afez / loadz debug pages (landing in a later PR). - `session_snapshot.go` — the value-typed snapshot record the debug surface consumes; no live locks escape. **Session helpers added** (pool-facing additions to files already touched by prior PRs, kept minimal): - `Session.loops sync.WaitGroup` + `WaitGoroutines()` — pool teardown blocks on this so `readLoop` / `heartbeatLoop` and their `notifyClosed → recordClose` callback chains fully unwind before `Close` returns. Prevents session goroutines from racing metric-var writes across test boundaries. - `Session.closeErr atomic.Pointer[error]` + `setCloseErr` / `closeError` — preserves the raw `Recv` error handed to `handleClose`. Pool surfaces this on consecutive-failure breaker trips so operators see the underlying server rejection (e.g. `FailedPrecondition` when the resource is still being created) instead of only the sentinel. **Supporting additions to existing files:** - `afe_picker.go` — const `defaultAfeRandomSubsetSize = 2` (power-of-two-choices K-choice default; matches Java). - `debug_tracer.go` — three new tag constants: `tagSessionPoolCreatePanic`, `tagSessionPoolConsecutiveFailuresTripped`, `tagSessionPoolCheckoutFailedCINil`. - `connpool.go` — `ChannelPickHintInto(ctx, *atomic.Int32)` context helper. No-op when the channel pool doesn't consume the hint. ## What does NOT land yet - `SessionPool` / `Invoker` interfaces (follow-up PR alongside sessionClient / sessionTable). - `bigtable.Client` integration (PR-3+). - Debug pages under `bigtable/debugview/` (later PR). ## Test plan - [x] `go build ./...` passes. - [x] `go vet ./internal/transport/` clean. - [x] `go test ./internal/transport/ -race -count=1 -short -skip 'AfeLbSim' -timeout=180s` — passes (32s wall). ~2200 LOC of new tests across pool lifecycle, scaling, consecutive-failure breaker, AFE integration, debug surface, snapshot rendering, plus a K-choice bench. --- # Reviewer guide ## Guide 1 — mutianf (human) ### What this PR does Adds `SessionPoolImpl`, the layer that sits above the per-AFE `sessionList` shipped in #20224 and consumes it via a two-tier picker (AFE first, then a ready session in that AFE). It owns the session lifecycle (open / active / closing / close hooks), server-driven scaling via `PoolSizer`, a consecutive-failure circuit breaker, and the debug/observability surface (histograms + ring buffers) that feeds sessionz/loadz. New files: 5 source, 6 test, ~4.9k LOC. Nothing outside `session_pool*.go` / `session_snapshot*.go` is new logic — the small edits elsewhere are hook-plumbing scaffolding already vetted by the session/AFE subagent reviewers. ### Recommended read order 1. **`session_pool.go`** — start here. Struct field layout with per-field ownership comments (`:104-179`), the `waiter` FIFO shape (`:94-101`), `CheckoutSession` two-tier pick + parking (`:235-310`), `Invoke` (`:465-559`), `Stats` (`:361-397`), `UpdateConfig` (`:402-431`), `pickerFromLoadBalancing` (`:439-461`). Skim `session_pool_test.go` (28 tests) — the FIFO waiter, Stats, and UpdateConfig behaviors are all covered there. 2. **`session_pool_lifecycle.go`** — hooks (`onActive:255`, `onClosing:308`, `onClose:336`), `recordSessionClose` once-CAS on `Session.poolCloseRecorded` (`:117-130`), `Close`'s 6-phase teardown (`:154-247`), `noteAbnormalCloseIfAny` breaker (`:363-392`), the three ticker loops (`:426-538`). Skim `session_pool_lifecycle_test.go` — every hook + `Close`. 3. **`session_pool_scaling.go`** — `Tick` (`:81-162`), `createSession` worker (`:164-274`), `scalingReason` (`:278-299`), `noDeadlineButCancellableContext` (`:301-311`). Skim `session_pool_scaling_test.go` — the `scalingInProgress` gate and panic-safety are the only non-obvious contracts. 4. **`session_pool_debug.go`** — `poolMetrics` (`:36-72`), `latencyHist` log2 histogram (`:160-228`), the four ring buffers (slow-vRPC, time-series, lifetimes, pick-history), `recordPickDecision` (`:366-387`). Skim `session_pool_debug_test.go` — mostly ring-cap and rate-computation coverage. 5. **`session_snapshot.go`** — mostly type defs. Focus on `PoolSnapshot` (`:452-594`) and `LoadBalancingSnapshot` (`:414-436`) as the debug-view contract. 6. **`session_pool_consecutive_failures_test.go`** and **`session_pool_afe_test.go`** — end-to-end behavior verification; useful for confirming intent. ### Flow of events - **CheckoutSession → Invoke → release.** `CheckoutSession` (`session_pool.go:235`) opportunistically kicks Tick if `sl.ReadyCount()==0`, snapshots the picker under `p.mu`, then two-tier picks outside the lock: `ReadyAfes()` → `PickAfe` → `Checkout(afeID)` (`:259-268`). Miss → park in the FIFO waiter queue (`:286-289`), bracket `waitersCount` for the sizer (`:291,300`). `Invoke` (`:465`) checks out, runs `sh.session.Invoke`, records latencies (`:508-523`), logs a slow-vRPC row if over threshold (`:524-557`); the deferred `sh.DecOutstanding()` + `noteVRpcOutcome` (`:493-496`) hands the OK-gated latency to the per-AFE PeakEwma tracker. Session release itself is driven by `OnSlotDrained` (installed at `session_pool_scaling.go:228-231`), which returns the handle to `sessionList` and calls `signalFree` — separate from the `defer` in `Invoke`. - **Background Tick.** `startTickLoop` (`session_pool_lifecycle.go:426`) fires every 1 s → `tickOnce` debounces via `tickPending` CAS (`:447-458`) → `Tick` (`session_pool_scaling.go:81`) samples uptimes, gates on `scalingInProgress`, calls `sizer.Decide()`, and on a positive delta reserves `pendingStarts += delta` + `spawns.Add(delta)` under `p.mu` (`:131-138`) then fans out one goroutine per session. Each `createSession` acquires the budget outside `p.mu`, dials via `streamFactory`, transfers `pendingStarts → startingSessions` in one lock (`:246-249`), starts the session, and blocks on `WaitGoroutines` so it stays on `p.spawns` until the session dies. - **Abnormal close → breaker trip.** `onClose` (`session_pool_lifecycle.go:336`) CAS's `closeRecorded`, calls `noteAbnormalCloseIfAny` (`:363`), which bumps `consecutiveFailures` and stores the raw error into `lastAbnormalCloseErr`. Crossing the threshold snapshots the poison, CAS-resets the counter, and calls `drainWaitersWithErr` — waiters get `*consecutiveFailureError` wrapping the last cause (so `errors.Is(err, ErrConsecutiveFailures)` and `status.Code(err)` both still work, `:60-82`). Counter only resets in `onActive` (`:292-293`) — a successful open, not a healthy vRPC. ### Key invariants 1. **Two-tier pick, no re-entrant `p.mu`.** `CheckoutSession` reads `p.picker` under `p.mu` (`session_pool.go:249-255`) then unlocks before calling picker/sessionList. `recordPickDecision` takes `pickerName` as a **parameter** (`session_pool_debug.go:366`, `session_pool.go:260-262`) precisely because the caller already holds no lock — but any new pool method that reads `p.picker.Name()` from a hot path must not re-take `p.mu`. 2. **Waiter FIFO with `waitersCount` bracketed.** Every `PushBack` bumps `waitersCount` (`session_pool.go:291`); every wake path (`ctx.Done`, `w.ready`) decrements it (`:294,300`). `removeWaiter` (`:316`) is idempotent via `w.elem != nil`; `signalFree` and `drainWaitersWithErr` nil out `elem` under `waitersMu` (`:329-358`). `Stats().PendingCount` reads `waitersCount.Load()` — this is the sizer's queue-depth input. 3. **Close-exactly-once accounting.** `sessionsClosed` and `closesByReason` bumps are gated by `Session.poolCloseRecorded.CompareAndSwap(false, true)` inside `recordSessionClose` (`session_pool_lifecycle.go:117-130`). `sh.closingRecorded` and `sh.closeRecorded` are per-handle CAS's protecting the lifetime histogram + the `OnClose` branch. `Close`'s Phase 1 pre-flips both CAS's on every handle (`:187-193`) so a concurrent mid-flight onClosing can't double-count. 4. **Breaker resets only on `onActive`.** `consecutiveFailures.Store(0)` and `lastAbnormalCloseErr.Store(nil)` live at `session_pool_lifecycle.go:292-293`. Not on per-vRPC OK — otherwise one long-lived healthy session would mask a run of failed opens. 5. **Hot path is atomics/RLocks; debug views take snapshots.** `Stats` is the only per-request path that briefly takes `p.mu` (`session_pool.go:362`); everything else on the vRPC path is atomic. Debug snapshotters copy under lock and format after release (`session_snapshot.go:452-594`). ### What NOT to worry about - **Session / vRPC layer itself** — shipped in #20213 / #20215 (state machine, one-in-flight, PeerInfo timing, retry oracle, heartbeat). - **Per-AFE `sessionList` I1-I6** — shipped in #20224, has its own tests. - **`PoolSizer` scaling formula** — already upstream (`pool_sizer.go`); this PR only wires it and consumes `ScaleDecision`. - **AFE pickers (`SimpleAfePicker` / `LeastInFlight` / `LeastLatency`)** — already upstream (`afe_picker.go`); this PR only builds them via `pickerFromLoadBalancing`. - **`SessionThrottler` / `AdaptiveSessionThrottler`** — already upstream; this PR consumes `Acquire` / `Release` / `UpdateConfig`. - **`ClientConfigurationManager` polling** — this pool receives `UpdateConfig` calls; the polling itself is elsewhere. ### Danger zones - **Re-entrant `p.mu` on picker access.** `recordPickDecision` intentionally takes `pickerName` as a param (`session_pool_debug.go:366`). Adding a new pool method that reads `p.picker.Name()` from within a `CheckoutSession` code path is a re-entrant deadlock; pass the name in or snapshot up-front. - **`startingSessions` / `pendingStarts` accounting.** Tick reserves `pendingStarts` under `p.mu` (`session_pool_scaling.go:131-138`), `createSession`'s `reserved` defer releases it on any early return (`:172-179`), and the transfer at `:246-249` is atomic under `p.mu`. `onActive` deletes from `startingSessions` (`session_pool_lifecycle.go:265`). Any new failure branch in `createSession` must preserve the invariant `pendingStarts + len(startingSessions) + Ready = "in-flight scale-up capacity"`. - **`budget.Acquire` blocks; must run OUTSIDE `p.mu`.** Currently at `session_pool_scaling.go:181`, deliberately after the `defer reserved` block and before any `p.mu.Lock()`. Moving it under the lock deadlocks scale-up under budget exhaustion. - **Slow-vRPC is fire-and-forget.** `recordSlowVRpc` (`session_pool_debug.go:301`) appends into a mutex-guarded ring; it's on the vRPC return path but bounded by threshold gating and one small mutex. Do not add I/O, logging fanout, or channel sends here. - **`newTestPool` / bootstrap defaults.** Test helpers construct the pool without a `ClientConfigurationManager`, so the bootstrap defaults from `defaultPoolConfig()` (`session_pool.go:218-226`) are the only config those tests see. Real callers always get `UpdateConfig` synchronously on registration — verify by reading `NewSessionPoolImpl` end-to-end, not by trusting the bootstrap values. - **`Close` phase ordering.** Phase 4 (`poolCancel`) runs AFTER Phase 3 (`wg.Wait` on graceful closes) because Phase 2's `closeCtx` derives from `p.poolCtx`; swapping order strands in-flight graceful closes on a cancelled ctx. Phase 5 (`p.spawns.Wait`) blocks on every createSession goroutine's `WaitGoroutines` — that's why createSession blocks on `s.WaitGoroutines()` at `session_pool_scaling.go:272`. --- ## Guide 2 — mutianf-bot (automated reviewer) ### REAL HAZARDS to flag - **Re-entrant `p.mu` in pool methods called from `CheckoutSession`.** Anchor: `session_pool.go:235-310`. `p.mu` is dropped at `:255` before `PickAfe` / `Checkout` / `recordPickDecision` fire. Flag any newly-added helper called from that block that re-acquires `p.mu`, or any new method that reads `p.picker.Name()` without taking the name as a parameter (see the intentional parameter pattern at `session_pool_debug.go:366`). - **`budget.Acquire` under `p.mu`.** Currently correctly outside the lock at `session_pool_scaling.go:181`. `SessionThrottler.Acquire` blocks on the budget semaphore; calling it while holding `p.mu` would deadlock scale-up. Flag any code path that acquires `p.mu` before line `:181` or moves `Acquire` inside a `Lock`/`Unlock` bracket. - **`sync.Map` allocations on hit paths.** `bumpCloseReason` uses `Load` first, `LoadOrStore(k, new(atomic.Int64))` only on miss (`session_pool_lifecycle.go:102-111`) — this is the correct pattern. Flag any new `sync.Map.LoadOrStore(key, new(...))` call on a hot path that isn't gated by a preceding `Load` — that allocates on every hit. - **Waiter counter drift.** `waitersCount.Add(+1)` at `session_pool.go:291`, `Add(-1)` on both the `ctx.Done` branch (`:294`) and the `w.ready` branch (`:300`). Flag any new wake path, timeout branch, or early-return between `:291` and `:308` that doesn't decrement, and any new enqueue site that doesn't increment. Drift here corrupts the sizer's `PendingCount` input. - **Unbalanced `pendingStarts` / `startingSessions`.** Tick increments `pendingStarts` under `p.mu` at `session_pool_scaling.go:131-138`; `createSession`'s `reserved` defer at `:172-179` releases on early return; the transfer to `startingSessions` at `:246-249` is atomic; `onActive` deletes at `session_pool_lifecycle.go:265`; failed-start deletes at `session_pool_scaling.go:253-255`. Flag any new failure branch in `createSession` that returns without either the `reserved` defer or an explicit transfer/cleanup. - **Missing CAS on close-once flags.** `sessionsClosed.Add(1)` sits downstream of `s.poolCloseRecorded.CompareAndSwap(false, true)` at `session_pool_lifecycle.go:121`. Lifetime histogram sits downstream of `sh.closingRecorded.CompareAndSwap(false, true)` at `:187` and `:317`. `recordSessionClose` re-entry is dedup'd at `:190` and `:347`. Flag any new counter bump / histogram record that isn't gated by the corresponding CAS. - **Uncontrolled goroutine spawns.** Every pool-spawned goroutine must be tracked so `Close` phases 3/5/6 can drain: createSession workers use `p.spawns.Add`/`Done` under `p.mu` (`session_pool_scaling.go:137-138,149`); `spawnTickOnce` uses the same pattern (`session_pool_lifecycle.go:464-476`); the mid-Close onActive dispatch adds at `:274`. Session-owned goroutines (`readLoop`/`heartbeatLoop`) are tracked on `Session.loops` and drained via `WaitGoroutines`. Flag any bare `go func(){…}()` in these files with no WaitGroup entry. - **Hot-path lock churn.** `p.mu` is briefly held on the fast path only at `session_pool.go:239-241` and `:249-255`. `Stats` takes `p.mu` (`:362`). Flag any new `p.mu.Lock()` inserted into `Invoke` (`:465-559`) or the pre-park stretch of `CheckoutSession` outside the two existing points. ### FALSE-POSITIVE SUPPRESSION - **"Package name should be `transport`, not `internal`."** The directory is `bigtable/internal/transport/` but the `package` declaration is `package internal` (confirmed at `session_pool.go:20`). Package name is correct as-is; do NOT flag. - **Local variable shadowing an exported type** where the local name is idiomatic (e.g. `afeID` local vs `AfeID` type in `CheckoutSession`). Already resolved intentionally; do NOT re-raise variants. - **`goimports` / `gofmt` / column-alignment / trailing-newline nits.** CI (`goimports -l`, `gofmt -l`, `go vet`) already gates these. Bot echo is noise. - **Comments referencing PR #20213 / #20215 / #20224.** Stacked-PR context, not stale references; do NOT suggest removal. - **Test coverage complaints for `pool_sizer.go`, `afe_picker.go`, `session_list.go`, `session.go`, `session_vrpc.go`, `session_throttler.go`, `client_configuration_manager.go`, `default_client_config.go`.** All shipped in earlier PRs (#20213, #20215, #20224) with their own tests; out of scope here. - **"Missing error wrapping"** on internal-only calls where the caller already annotates via `fmt.Errorf("POOL %s ...: %w", ...)` or via `btopt.Debugf`. Do NOT suggest adding a second wrap. - **Retry loop / context propagation questions on `Session.Invoke`.** That's the Session layer (`session_vrpc.go`), out of scope for this PR. - **"Consider using `sync.RWMutex` instead of `sync.Mutex` on `p.mu`."** The pool holds `p.mu` for tens of nanoseconds at a time and never for read-heavy loops; the added atomic on `RLock`/`RUnlock` would cost more than it saves. Do NOT suggest. - **"Consider extracting anonymous goroutine into named function."** Style-only; do NOT suggest for the three ticker loops or the createSession worker. ### SCOPE BOUNDARY Comment ONLY on: - `bigtable/internal/transport/session_pool.go` - `bigtable/internal/transport/session_pool_lifecycle.go` - `bigtable/internal/transport/session_pool_scaling.go` - `bigtable/internal/transport/session_pool_debug.go` - `bigtable/internal/transport/session_snapshot.go` - `bigtable/internal/transport/session_pool_*_test.go` - `bigtable/internal/transport/session_snapshot_test.go` Do NOT comment on additions to: - `session.go` / `session_vrpc.go` (WaitGoroutines / closeError additions — vetted) - `connpool.go` (`ChannelPickHintInto` helper — vetted) - `afe_picker.go` (`defaultAfeRandomSubsetSize` constant — vetted) - `debug_tracer.go` (3 new tags — vetted) These are supporting scaffolding, already reviewed by the 3 subagent reviewers in this stack. Only re-raise if something looks actively unsafe. ### EFFORT SCALING - ~4.9k LOC across 12 files. Do NOT paginate uniformly. - **First pass — the 4 hot source files, in this order:** 1. `session_pool.go` (559 LOC) 2. `session_pool_lifecycle.go` (538 LOC) 3. `session_pool_scaling.go` (311 LOC) 4. `session_pool_debug.go` (416 LOC) - **Second pass ONLY if a first-pass finding needs corroboration:** `session_snapshot.go` (594 LOC, mostly type defs), and the tests. Tests use `newTestPool`, which skips config wiring — do NOT flag bootstrap defaults on tests as if they were production paths. - If a first-pass finding is a real hazard from the list above, cite the file:line and the exact anchor pattern it violates. Do not file speculative "consider" comments.Sushan Bhattarai · 683eda8c · 2026-07-28
- 4.1ETVfeat(storage): merge support for bi-directional multiple range reads (#11377) * feat(storage): merge support for bi-directional multiple range reads * add license header * get object metadata for regular grpc reads * docs change * fix offset and remaining bytes calculation * fix(storage): use mutex for accessing concurrent vars * remove TestOpenAppendableWriterUnsupportedEmulated * Add back retry conformance tests to json * retry tests assume writes are conditionally idempoten --------- Co-authored-by: Daniel B <danielduhh@gmail.com>Brenna N Epp · b4d86a52 · 2025-01-08
- 3.9ETVfeat(auditmanager): add new clients (#13842) PiperOrigin-RevId: 861758611Alex Hong · 588fe0c8 · 2026-02-12
- 3.5ETVfeat(iam/admin): regenerate client (#11570) This client has some handwritten code in it that causes failures when generating it normally. I generated the client and hand patched around handwritten code as to not create a breaking change. Changes made: Keep old variants of RPCs that did not generate an iterator helper to avoid breaking changes: ListRoles, QueryGrantableRoles, and QueryTestablePermissions Add variants of above RPCs that do return an iterator: ListRolesIter, QueryGrantableRolesIter, and QueryTestablePermissionsIter Preserve API for GetIamPolicy and SetIamPolicy. Fixes: #8219Cody Oss · eab87d73 · 2025-02-07
- 3.5ETVfeat(hypercomputecluster): add new client (#13423) PiperOrigin-RevId: 823382075Alex Hong · 945efa94 · 2025-12-02
- 3.4ETVchore: onboard a new mapmanagement library (#14467)Noah Dietz · 26125607 · 2026-04-22
- 3.3ETVchore: onboard new library appoptimize (#14393) Onboard new API client for `google/cloud/appoptimize/v1beta` Internal bug b/500129743Noah Dietz · e247d48b · 2026-04-08
- 3.2ETVfeat(bigtable): add Session lifecycle (Start, Close, ForceClose, readLoop, heartBeatLoop) (#20215) ## Summary Third and final PR in the Session core stack. Adds the lifecycle orchestration that drives a Session from \`Start\` through teardown, plus the \`readLoop\` that dispatches server frames to the vRPC handlers landed in #20213. **Stacks on** #20213 (Session vRPC dispatch + slot lifecycle). Deletes the minimal \`ForceClose\` stub that PR shipped, replacing it with the full lifecycle-shaped body. ### What lands **New file: \`session_lifecycle.go\`** (~640 LOC) - \`Session.Start(ctx, OpenSessionRequest)\` — transitions New→Starting, Sends the OpenSession frame, fires \`onStart\`, spawns \`readLoop\` + \`heartBeatLoop\`. Wraps a failed \`Send\` as \`codes.Unavailable\` so retry plumbing treats pre-wire OpenSession loss the same as any other transport-side loss. - \`ForceClose\` — full body: \`transitionTo(Closed)\` → \`setCloseReason\` → \`notifyClosing\` (once) → \`cancelActiveRPCs\` → \`signalQuiescent\` → \`notifyClosed\`. - \`Close(ctx, CloseSessionRequest)\` — graceful drain: Ready→Closing, waits on \`quiescent\`, sends CloseSession, transitions to WaitServerClose, arms the pool's stuck-session sweep to eventually ForceClose if the server never confirms. - \`notifyClosing\` / \`notifyClosed\` — once-guarded hook dispatchers with the strict \`onClosing\`-precedes-\`onClose\` ordering enforced. - \`readLoop(ctx)\` — Recv-loop that dispatches every SessionResponse variant via \`handleSessionResponse\`, drives \`handleClose\` on stream termination, records \`msgsRecv\` counters + resets heartbeat deadline on every recognized frame. - \`handleSessionResponse\` — dispatch switch to \`handleOpenSession\` / \`handleVRPCResponse\` / \`handleErrorResponse\` / \`handleSessionParameters\` / \`handleGoAway\` / \`handleSessionRefreshConfig\`. Heartbeat frames reset the deadline; unknown-payload branch records a debug tag but doesn't reset (so a misbehaving server can't keep the watchdog satisfied with junk). - \`handleOpenSession\` — Starting→Ready transition, PeerInfo extract from the bidi header (synchronous, matches Java's onHeaders synchrony), fires \`onActive\`. - \`handleGoAway\` / \`handleClose\` / \`handleErrorResponse\` / \`handleSessionParameters\` / \`handleSessionRefreshConfig\` — protocol-level handlers. - \`heartBeatLoop\` — Timer + \`heartbeatWake\` reactive-wake pair. Only enforces the deadline while a vRPC is in flight (idle sessions legitimately receive no server heartbeats). See SESSION_SPEC.md #7. - \`peerInfoExtracter\` — parses the \`bigtable-peer-info\` header, stamps \`s.peerInfo\` atomically before \`onActive\` fires. - \`closeReasonLabel\` / \`closeReasonToCause\` — CloseSessionRequest.Reason → string / sentinel error mapping for close-reason attribution. **New file: \`session_lifecycle_test.go\`** (~660 LOC) — 30+ tests covering Start, ForceClose, Close, readLoop, handleSessionResponse dispatch, handleOpenSession + peerInfo parsing, handleGoAway, handleClose, handleErrorResponse (rpc_id=0 harmlessly drops via routeVRPCFrame guards), heartBeatLoop (reactive wake, idle-gate, missed-heartbeat ForceClose). **Edits to \`session_vrpc.go\`** - Delete the minimal \`ForceClose\` stub introduced in the prior PR — the full body lives in \`session_lifecycle.go\` now. **Edits to \`session_test.go\`** - \`hookCounts\` helper (start/active/close callback counters). - \`setSlotForTest\` helper (seed the in-flight slot from tests). ### Stack 1. #20211 — Session debug surface 2. #20213 — Session vRPC dispatch + slot lifecycle 3. **This PR** — Session lifecycle ## Test plan - [x] \`go build ./internal/transport/\` passes - [x] \`go vet ./internal/transport/\` clean - [x] \`go test ./internal/transport/ -count=1 -short -timeout 90s\` passes — 30+ new lifecycle tests plus all pre-existing testsSushan Bhattarai · b9e53c62 · 2026-07-27
- 3.1ETVfeat(bigtable): add SessionClient + SessionTable + lazyPool (#20228) ## Summary Third of five PRs porting the session-pool infrastructure from `feat/bigtable-sessionz-debug` into upstream. Introduces `internal/session/` — a proto-native SessionClient + SessionTable API sitting on top of PR-2's SessionPoolImpl. - `internal/session/api.go` — public interfaces (`ChannelPool`, `Config`, `SessionClient`, `SessionTableAPI`, `DebugAccess`). - `internal/session/client.go` — `SessionClient` impl: dedicated channel pool (no primer), `ClientConfigurationManager` wiring, `OpenSessionTable` / `OpenAuthorizedView` / `OpenMaterializedView` factories that mint lazily-opened per-resource pools keyed by `{resource, permission}`. - `internal/session/table.go` — `SessionTable` impl. Two `*lazyPool` (read + write); MV is read-only (write pool nil, `MutateRow` returns `ErrWriteNotSupported`). `stampAttempt` sources per-attempt `cluster_id` / `zone_id` / peer fields from typed `InvokeResult.ClusterInfo` and `InvokeResult.PeerInfo` per CLIENT_SIDE_METRICS_SPEC #1. - `internal/session/lazy_pool.go` — `Invoker` + `SessionPool` interfaces + open-on-first-use lazy wrapper. Failed opens are NOT cached; the next call retries. - `internal/session/debug.go` — `DebugAccess` impl surfacing pool snapshots for sessionz/loadz/channelz/configz. ### Transport additions to support the above - `transport/debug_api.go` (new) — `SessionDebugProvider` / `ChannelDebugProvider` / `ConfigDebugProvider` interfaces + `ChannelPoolDebug` / `SessionRef` DTOs. Lives in transport (not bigtable) so `bigtable.Client` and `internal/session.SessionClient` can implement without an import cycle. - `transport/diverter.go` — `sessionPicks` / `classicPicks` counters, `DiverterSnapshot`, `Snapshot()`. - `transport/connpool.go` — `ChannelSnapshot` + `ChannelPoolSnapshot` type + method, `WithInstanceName` / `WithAppProfile` options for channelz labelling. - `transport/debug_tracer.go` — exported `DebugTag` + `RecordDebugTag` + `TagSessionAttemptNilClusterInfo` / `TagSessionAttemptEmptyClusterID` catalog constants. - `transport/session_descriptors.go` — `SessionType.ProtoName()` for human-readable pool identifiers. - `transport/direct_access_checker.go` — renames `newPingAndWarmDirectAccessChecker` → `NewPingAndWarmDirectAccessChecker` (constructor exported) and nil-guards `primer.Prime` so session-based clients can pass a nil primer. Session clients warm channels on-demand via `OpenSession`, not eagerly at pool-init. ### Lifecycle correctness `sessionClient.Close()` snapshots owned resources under `poolsMu` and releases the lock before running `Close` / `Shutdown` / `Cancel` calls, so a snapshot method holding `poolsMu` never deadlocks teardown. Post-Close Opens surface a distinct `ErrSessionClientClosed` sentinel rather than misleading `errReadPoolNil` / `ErrWriteNotSupported`. ## Stack - **PR-1** #20224 (sessionList) — merged - **PR-2** #20225 (SessionPoolImpl) — open; this PR stacks on it - **PR-4** (this PR) - PR-5 will follow with the bigtable-package integration + debugview. The diff on this PR includes PR-2's commits until #20225 merges into main. ## Test plan - [x] `go build ./internal/session/ ./internal/transport/ ./...` - [x] `go test ./internal/session/ -race -count=1 -short -timeout=180s` - [x] `go test ./internal/transport/ -race -count=1 -short -skip 'AfeLbSim|TestHighQpsSession' -timeout=240s` - [x] `gofmt -l ./internal/session/ ./internal/transport/` — clean - [x] `go vet ./internal/session/ ./internal/transport/` — clean - [x] Reviewed against the 4 behavioral specs (SESSION_SPEC, SESSION_CLIENT_SPEC, SESSION_POOL_SPEC, CLIENT_SIDE_METRICS_SPEC) and SESSION_COMPONENT_SPEC (boundary rules) — PASS.Sushan Bhattarai · ab2c96c3 · 2026-07-28
- 3.1ETVfeat(firestore): refactor pipeline API for uniform stage options (#14322) Design decision [here](https://docs.google.com/document/d/1diPUhXr2kZh69KcAsuOMZMj_rfcykGsE2-phmBIt8vs/edit?usp=sharing) **Overview** This PR standardizes the API surface for the Firestore Pipeline builder in the Go SDK, ensuring idiomatic conformance and aligning closely with cross-platform design paradigms (Java/Node.js). It replaces complex stage-specific configuration structs with a universal functional option pattern and introduces typed slice constants for varied expressions. **Key Features and API Changes:** * **Universal `RawOptions` Escape Hatch:** Introduced `RawOptions` to serve as a universal pass-through for backend options that might not be natively modeled by the Go SDK (e.g., `RawOptions{"limit": 5, "distance_field": "dist"}`). This deprecates intermediate parameter structs like `PipelineFindNearestOptions` across pipeline stages. * **Standardized Variadic Arguments with Slices:** Methods that formerly relied on raw variadic argument unpacking (e.g. `...Ordering`, `...any`) have been updated to accept typed slices. To preserve clean query ergonomics, we introduced domain-specific helper constructors (`Orders()`, `Fields()`, `Accumulators()`, `Selectables()`). * `Pipeline.Sort(Orders(Ascending(FieldOf("f")), Descending(FieldOf("__name__"))))` * `Pipeline.Aggregate(Accumulators(Count(DocumentID).As("total")))` * `Pipeline.Select(Fields("a", "b"))` * **Functional Stage Options Migration:** Implemented standard functional options for specialized stage modifiers: * `WithFindNearestLimit()`, `WithFindNearestDistanceField()` * `WithAggregateGroups()` * `WithUnnestIndexField()` * `WithExplainMode()` **Internal Code Simplification & Testing:** * **Flattened Stage Hierarchy (`baseStage` removal):** Removed the `baseStage` struct embedding across our internal pipeline stages. We replaced it with explicit `options map[string]any` fields within each stage and introduced a shared `stageOptionsToProto(options)` helper. This simplifies the internal struct hierarchy, removes opaque state inherited from struct embedding, and makes configuration serialization significantly more explicit and easier to trace inside each stage's `toProto()` implementation. * **Pipeline Stage Proto Parsing:** Refactored internal `pipelineStage` structures (`newAggregateStage`, `newFindNearestStage`, `newSortStage`, etc.) to reliably map the new `RawOptions` down to standard `*pb.Value` payloads. * **Tests for `AlwaysUseImplicitOrderBy`:** Added dedicated test cases (`TestQuery_AlwaysUseImplicitOrderBy`) to explicitly verify the behavior of the existing `WithAlwaysUseImplicitOrderBy` client toggle, ensuring it appropriately serializes implicit inequality fields and `__name__` into the `OrderBy` proto requests. * **Integration Suite Migration:** Fully migrated the integration suite inside `pipeline_integration_test.go` and `pipeline_test.go` to assert against the newly stabilized `RawOptions` and slice-builder interfaces. --------- Co-authored-by: Alex Hong <9397363+hongalex@users.noreply.github.com>Baha Aiman · ca7c3699 · 2026-04-06
- 2.8ETVrefactor(bigtable): extract metrics tracer into internal/metrics for classic+session reuse (#20099) Summary 1. Decouple the metrics Tracer from gax invoker 2. Move metrics into internal package 3. Keep NoopMetricsTracer for aliasing.Sushan Bhattarai · 2b621805 · 2026-07-08
- 2.7ETVfeat: update API sources and regenerate (#14621) This change is from "Update & generate" for Go in the playbook.Tomo Suzuki · 6641db88 · 2026-05-20
- 2.4ETVfeat(spanner): add dynamic channel pool (#14611) Split of https://github.com/googleapis/google-cloud-go/pull/14604 Internal reference: go/go-dcp-designrahul2393 · 51a53cee · 2026-05-26
- 2.4ETVchore(spanner): integrate location aware routing with RPCs (#13877) - Add `locationAwareSpannerClient` wrapper that intercepts RPCs and routes them to the server endpoint resolved by the existing `channelFinder`, falling back to the default gRPC channel when no endpoint is available - Add `endpointClientCache` that creates and caches per-address gRPC connections - Add transaction affinity tracking so Commit/Rollback route to the same server that handled the transaction, with read-only transactions routed independently per-request based on key ranges - Move request preparation (`prepareReadRequest`, etc.) and response observation (`observePartialResultSet`, etc.) from transaction/batch code into the client wrapper, keeping routing concerns in one placerahul2393 · b42935db · 2026-03-10
- 2.3ETVfeat(firestore): Introduce new functions and literals stage (#14114) [go/firestore-query-tracker](http://go/firestore-query-tracker) : Introduce all the functions marked "P" (planned) in Go column O and "C" (Complete) in Backend column H in "Firestore Features (Pipeline)" sheet. The sheet above specified only the function name. The args for all the functions were determined from backend implementation in - google3/java/com/google/cloud/datastore/client/firestorev1/Functions.java - google3/java/com/google/cloud/datastore/client/firestorev1/Stages.java The documentation is only available for the functions marked "C" (Complete) in Docs column I. The same documentation has added to the client in this PR. - google3/googledata/devsite/_common/en/_shared/firestore/pipeline/functions/_shared/Baha Aiman · 7724d79f · 2026-03-04
- 2.2ETVfeat(gkerecommender): add new clients (#13238) PiperOrigin-RevId: 805881130Chris Smith · 8f604ff6 · 2025-10-27
- 2.1ETVfeat(apiregistry): add new clients (#13525) PiperOrigin-RevId: 848064295shollyman · 1c4f4008 · 2025-12-30
- 2.0ETVchore(internal/librariangen): add librariangen (#12614) * Add internal/librariangen to go.work. * Expect source_roots instead of source_paths in generate-request.json for now. closes: googleapis/librarian#791Chris Smith · 420edfe5 · 2025-07-30