Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ labels, aggregation rules, PromQL examples, and admission metric migration.
| `duckgres_session_admission_reclaim_reservation_capacity` | Gauge | Cleanup-ownership slot capacity for this control-plane process (4096 per reclaimer by default) |
| `duckgres_session_admission_reclaim_reservation_rejections_total{reason}` | Counter | Reservations rejected because capacity was `full`, the reclaimer was `closed`, or the exact reference was a `duplicate` |
| `duckgres_session_start_duration_seconds{org,protocol,outcome}` | Histogram | Authenticated PostgreSQL session bootstrap through flushed `ReadyForQuery` |
| `duckgres_postgres_session_start_total{org,outcome,reason}` | Counter | Exactly one terminal result per authenticated PostgreSQL session start after server retries; `outcome` is `success\|failure` and bounded reasons distinguish operator-actionable failures from client/lifecycle noise |
| `duckgres_flight_rpc_duration_seconds{method}` | Histogram | Flight ingress RPC duration by method |
| `duckgres_flight_ingress_sessions_total{outcome}` | Counter | Flight ingress session outcomes (`created|reused|auth_failed|rate_limited|create_failed|token_invalid`) |
| `duckgres_flight_sessions_reaped_total{trigger}` | Counter | Number of Flight auth sessions reaped (`trigger=periodic|forced`) |
Expand Down
108 changes: 85 additions & 23 deletions controlplane/control.go
Original file line number Diff line number Diff line change
Expand Up @@ -1144,7 +1144,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
server.RecordSuccessfulAuthAttempt(cp.rateLimiter, remoteAddr)
clog.Info("User authenticated.")
sessionStart := observe.BeginSessionStart(orgID, "postgres")
defer sessionStart.Finish("error")
defer sessionStart.Finish("error", observe.SessionStartReasonUnknown)

// Resolve the requested worker shape from the connection-string startup
// options (duckgres.worker_cpu / worker_memory / worker_ttl), layered on
Expand All @@ -1159,6 +1159,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
}
workerProfile, profileWarns, orgProfileApplied, profileErr := cp.resolveWorkerProfile(startupOptions, orgProfileDefaults)
if profileErr != nil {
sessionStart.Finish("error", observe.SessionStartReasonClient)
clog.Warn("Rejected worker profile.", "error", profileErr)
_ = server.WriteErrorResponse(writer, "FATAL", "22023", profileErr.Error())
_ = writer.Flush()
Expand All @@ -1181,6 +1182,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
rawS3Cache, s3CacheRequested := startupOptions[server.S3CacheGUCName]
if s3CacheRequested {
if err := server.ValidateS3CacheOption(rawS3Cache); err != nil {
sessionStart.Finish("error", observe.SessionStartReasonClient)
clog.Warn("Rejected s3_cache startup option.", "error", err)
_ = server.WriteErrorResponse(writer, "FATAL", "22023", err.Error())
_ = writer.Flush()
Expand All @@ -1197,6 +1199,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
// hitting a partially-migrated catalog. The client gets a clear error
// and can retry after the migration completes.
if cp.orgRouter.IsMigratingForOrg(orgID) {
sessionStart.Finish("error", observe.SessionStartReasonLifecycle)
clog.Info("Connection rejected during DuckLake migration.")
_ = server.WriteErrorResponse(writer, "FATAL", "57P03",
"DuckLake catalog upgrade in progress for your organization, please retry in a few moments")
Expand All @@ -1212,6 +1215,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
// connection advisory lock.)
if whState, ok := cp.configStore.OrgWarehouseStatus(orgID); ok &&
whState == string(configstore.ManagedWarehouseStateResharding) {
sessionStart.Finish("error", observe.SessionStartReasonLifecycle)
clog.Info("Connection rejected during metadata-store reshard.")
_ = server.WriteErrorResponse(writer, "FATAL", "57P03",
"metadata-store reshard in progress for your organization, please retry shortly")
Expand All @@ -1226,6 +1230,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
// the stack absence is expected and the client should be told to
// retry rather than receive a misleading auth-style error.
whState, orgExists := cp.configStore.OrgWarehouseStatus(orgID)
sessionStart.Finish("error", missingOrgStackReason(whState, orgExists))
switch {
case !orgExists:
_ = server.WriteErrorResponse(writer, "FATAL", "28000", "no org configured for user")
Expand Down Expand Up @@ -1262,14 +1267,14 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
rebalancer = cp.rebalancer
}
if cp.isDraining() {
sessionStart.Finish("draining")
sessionStart.Finish("draining", observe.SessionStartReasonLifecycle)
_ = server.WriteErrorResponse(writer, "FATAL", "57P03", "control plane is draining, retry shortly")
_ = writer.Flush()
return
}
preReadyCtx, finishPreReady, err := cp.beginPreReadyHandshake(context.Background())
if err != nil {
sessionStart.Finish("draining")
sessionStart.Finish("draining", observe.SessionStartReasonLifecycle)
_ = server.WriteErrorResponse(writer, "FATAL", "57P03", "control plane is draining, retry shortly")
_ = writer.Flush()
return
Expand All @@ -1286,6 +1291,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
defer server.CancelClientConn(tmpCC)
server.SendInitialParams(tmpCC)
if err := writer.Flush(); err != nil {
sessionStart.Finish("error", observe.SessionStartReasonTransport)
clog.Error("Failed to flush initial params.", "error", err)
return
}
Expand Down Expand Up @@ -1330,7 +1336,8 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
sessions.DestroySession(createdPID)
executor = nil
}
sessionStart.Finish(controlPlaneSessionStartOutcome(err))
outcome, reason := controlPlaneSessionStartResult(err)
sessionStart.Finish(outcome, reason)
clog.Error("Failed to create session.", "error", err)
code, message := sessionCreationErrorResponse(err)
_ = server.WriteErrorResponse(writer, "FATAL", code, message)
Expand Down Expand Up @@ -1360,7 +1367,13 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
attachContextErr := attachCtx.Err()
attachCancel()
if probeErr != nil {
sessionStart.Finish(controlPlaneSessionStartOperationOutcome(probeErr, attachContextErr, cp.isDraining()))
outcome, reason := controlPlaneSessionStartOperationResult(
probeErr,
attachContextErr,
cp.isDraining(),
observe.SessionStartReasonMetadataStore,
)
sessionStart.Finish(outcome, reason)
clog.Error("Failed to detect attached catalogs.", "error", probeErr)
_ = server.WriteErrorResponse(writer, "FATAL", "XX000", "failed to detect attached catalogs")
_ = writer.Flush()
Expand All @@ -1371,6 +1384,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
var ok bool
effectiveCatalog, ok = resolveEffectiveCatalog(requestedCatalog, duckLakeAttached)
if !ok {
sessionStart.Finish("error", observe.SessionStartReasonMetadataStore)
clog.Warn("Postgres connection rejected: requested catalog is not available for this connection.",
"requested", requestedCatalog, "ducklake_attached", duckLakeAttached)
msg := "no catalog is available for this connection"
Expand Down Expand Up @@ -1413,7 +1427,13 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
if err := sessionmeta.InitSessionDatabaseMetadataWithAccess(initCtx, executor, effectiveCatalog, metadataAccess); err != nil {
initContextErr := initCtx.Err()
initCancel()
sessionStart.Finish(controlPlaneSessionStartOperationOutcome(err, initContextErr, cp.isDraining()))
outcome, reason := controlPlaneSessionStartOperationResult(
err,
initContextErr,
cp.isDraining(),
observe.SessionStartReasonMetadataStore,
)
sessionStart.Finish(outcome, reason)
clog.Error("Failed to initialize session database metadata.", "database", database, "error", err)
_ = server.WriteErrorResponse(writer, "FATAL", "XX000", "failed to initialize session database metadata")
_ = writer.Flush()
Expand All @@ -1433,7 +1453,13 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
spCancel()
if err != nil {
if source == sessionDefaultSourceConfiguredCatalog {
sessionStart.Finish(controlPlaneSessionStartOperationOutcome(err, spContextErr, cp.isDraining()))
outcome, reason := controlPlaneSessionStartOperationResult(
err,
spContextErr,
cp.isDraining(),
observe.SessionStartReasonMetadataStore,
)
sessionStart.Finish(outcome, reason)
clog.Error("Failed to apply session default catalog.", "catalog", effectiveCatalog, "error", err)
_ = server.WriteErrorResponse(writer, "FATAL", "XX000", "failed to apply default catalog")
_ = writer.Flush()
Expand All @@ -1456,7 +1482,13 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
initContextErr := initCtx.Err()
initCancel()
if err != nil {
sessionStart.Finish(controlPlaneSessionStartOperationOutcome(err, initContextErr, cp.isDraining()))
outcome, reason := controlPlaneSessionStartOperationResult(
err,
initContextErr,
cp.isDraining(),
observe.SessionStartReasonMetadataStore,
)
sessionStart.Finish(outcome, reason)
clog.Error("Failed to apply passthrough session default catalog.", "command", cmd, "error", err)
_ = server.WriteErrorResponse(writer, "FATAL", "XX000", "failed to apply default catalog")
_ = writer.Flush()
Expand Down Expand Up @@ -1521,6 +1553,7 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
// s3_cache=off must never silently run through the cache.
if s3CacheRequested {
if err := server.ApplyConnectionS3CacheOption(cc, rawS3Cache); err != nil {
sessionStart.Finish("error", observe.SessionStartReasonWorker)
clog.Error("Failed to apply s3_cache startup option.", "error", err)
_ = server.WriteErrorResponse(writer, "FATAL", "XX000", err.Error())
_ = writer.Flush()
Expand All @@ -1534,24 +1567,26 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
disconnect := disconnectWatcher.Stop()
drained := finishPreReady()
if disconnect.ClientCanceled {
sessionStart.Finish("canceled")
sessionStart.Finish("canceled", observe.SessionStartReasonCanceled)
return
}
if drained {
sessionStart.Finish("draining")
sessionStart.Finish("draining", observe.SessionStartReasonLifecycle)
return
}

// Send ReadyForQuery to signal that the handshake is complete
if err := server.WriteReadyForQuery(writer, 'I'); err != nil {
sessionStart.Finish("error", observe.SessionStartReasonTransport)
clog.Error("Failed to send ReadyForQuery.", "error", err)
return
}
if err := writer.Flush(); err != nil {
sessionStart.Finish("error", observe.SessionStartReasonTransport)
clog.Error("Failed to flush writer.", "error", err)
return
}
sessionStart.Finish("success")
sessionStart.Finish("success", observe.SessionStartReasonNone)
if orgID != "" {
observeOrgPgSessionAccepted(orgID, passthroughUser)
}
Expand All @@ -1564,33 +1599,60 @@ func (cp *ControlPlane) handleConnection(conn net.Conn) {
}
}

func controlPlaneSessionStartOutcome(err error) string {
func controlPlaneSessionStartResult(err error) (outcome, reason string) {
var capacityErr *WorkerCapacityExhaustedError
var rejection *configstore.OrgConnectionAdmissionRejectedError
switch {
case err == nil:
return "success"
case errors.As(err, &capacityErr), errors.As(err, &rejection):
return "capacity"
return "success", observe.SessionStartReasonNone
case errors.As(err, &capacityErr):
return "capacity", observe.SessionStartReasonCapacity
case errors.As(err, &rejection):
return "capacity", observe.SessionStartReasonClient
case errors.Is(err, context.Canceled):
return "canceled"
return "canceled", observe.SessionStartReasonCanceled
case errors.Is(err, context.DeadlineExceeded), errors.Is(err, ErrTooManyConnections):
return "timeout"
return "timeout", observe.SessionStartReasonCapacity
case errors.Is(err, ErrSessionManagerDraining):
return "draining"
return "draining", observe.SessionStartReasonLifecycle
default:
return "error"
return "error", observe.SessionStartReasonWorker
}
}

func controlPlaneSessionStartOperationOutcome(err, contextErr error, draining bool) string {
func controlPlaneSessionStartOperationResult(
err, contextErr error,
draining bool,
reason string,
) (outcome, classifiedReason string) {
if draining {
return "draining"
return "draining", observe.SessionStartReasonLifecycle
}
terminalErr := err
if contextErr != nil {
return controlPlaneSessionStartOutcome(contextErr)
terminalErr = contextErr
}
outcome, _ = controlPlaneSessionStartResult(terminalErr)
switch outcome {
case "success":
return outcome, observe.SessionStartReasonNone
case "canceled":
return outcome, observe.SessionStartReasonCanceled
case "draining":
return outcome, observe.SessionStartReasonLifecycle
default:
return outcome, reason
}
}

func missingOrgStackReason(warehouseState string, orgExists bool) string {
if !orgExists {
return observe.SessionStartReasonClient
}
if warehouseState == "" || warehouseState == string(configstore.ManagedWarehouseStateReady) {
return observe.SessionStartReasonControlPlane
}
return controlPlaneSessionStartOutcome(err)
return observe.SessionStartReasonLifecycle
}

func authorizedClientSearchPath(searchPath string, policy *server.QueryAccessPolicy) string {
Expand Down
91 changes: 68 additions & 23 deletions controlplane/control_cancel_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"github.com/posthog/duckgres/controlplane/configstore"
"github.com/posthog/duckgres/server"
"github.com/posthog/duckgres/server/flightclient"
"github.com/posthog/duckgres/server/observe"
)

func TestCreateSessionWithRegisteredCancel_CancelQueryCancelsWait(t *testing.T) {
Expand Down Expand Up @@ -196,52 +197,96 @@ func TestSessionCreationErrorResponse(t *testing.T) {
})
}

func TestControlPlaneSessionStartOutcome(t *testing.T) {
func TestControlPlaneSessionStartResult(t *testing.T) {
tests := []struct {
name string
err error
want string
name string
err error
wantOutcome string
wantReason string
}{
{name: "success", want: "success"},
{name: "canceled", err: context.Canceled, want: "canceled"},
{name: "deadline", err: context.DeadlineExceeded, want: "timeout"},
{name: "queue timeout", err: ErrTooManyConnections, want: "timeout"},
{name: "draining", err: ErrSessionManagerDraining, want: "draining"},
{name: "worker capacity", err: NewWorkerCapacityExhaustedError(time.Second), want: "capacity"},
{name: "success", wantOutcome: "success", wantReason: observe.SessionStartReasonNone},
{name: "canceled", err: context.Canceled, wantOutcome: "canceled", wantReason: observe.SessionStartReasonCanceled},
{name: "deadline", err: context.DeadlineExceeded, wantOutcome: "timeout", wantReason: observe.SessionStartReasonCapacity},
{name: "queue timeout", err: ErrTooManyConnections, wantOutcome: "timeout", wantReason: observe.SessionStartReasonCapacity},
{name: "draining", err: ErrSessionManagerDraining, wantOutcome: "draining", wantReason: observe.SessionStartReasonLifecycle},
{name: "worker capacity", err: NewWorkerCapacityExhaustedError(time.Second), wantOutcome: "capacity", wantReason: observe.SessionStartReasonCapacity},
{
name: "admission hard rejection",
err: &configstore.OrgConnectionAdmissionRejectedError{
Reason: configstore.OrgConnectionAdmissionRejectedOrgVCPU,
RequestedVCPUs: 4,
MaximumVCPUs: 2,
},
want: "capacity",
wantOutcome: "capacity",
wantReason: observe.SessionStartReasonClient,
},
{name: "generic error", err: errors.New("bootstrap failed"), want: "error"},
{name: "generic error", err: errors.New("bootstrap failed"), wantOutcome: "error", wantReason: observe.SessionStartReasonWorker},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := controlPlaneSessionStartOutcome(tt.err); got != tt.want {
t.Fatalf("controlPlaneSessionStartOutcome(%v) = %q, want %q", tt.err, got, tt.want)
gotOutcome, gotReason := controlPlaneSessionStartResult(tt.err)
if gotOutcome != tt.wantOutcome || gotReason != tt.wantReason {
t.Fatalf("controlPlaneSessionStartResult(%v) = (%q, %q), want (%q, %q)",
tt.err, gotOutcome, gotReason, tt.wantOutcome, tt.wantReason)
}
})
}
}

func TestControlPlaneSessionStartOperationOutcomePrefersContext(t *testing.T) {
func TestControlPlaneSessionStartOperationResultPrefersContext(t *testing.T) {
operationErr := errors.New("metadata initialization failed")
if got := controlPlaneSessionStartOperationOutcome(operationErr, context.Canceled, false); got != "canceled" {
t.Fatalf("canceled operation outcome = %q, want canceled", got)
tests := []struct {
name string
contextErr error
draining bool
wantOutcome string
wantReason string
}{
{name: "client canceled", contextErr: context.Canceled, wantOutcome: "canceled", wantReason: observe.SessionStartReasonCanceled},
{name: "operation timed out", contextErr: context.DeadlineExceeded, wantOutcome: "timeout", wantReason: observe.SessionStartReasonMetadataStore},
{name: "ordinary metadata error", wantOutcome: "error", wantReason: observe.SessionStartReasonMetadataStore},
{name: "draining wins", contextErr: context.Canceled, draining: true, wantOutcome: "draining", wantReason: observe.SessionStartReasonLifecycle},
}
if got := controlPlaneSessionStartOperationOutcome(operationErr, context.DeadlineExceeded, false); got != "timeout" {
t.Fatalf("timed-out operation outcome = %q, want timeout", got)

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
gotOutcome, gotReason := controlPlaneSessionStartOperationResult(
operationErr,
tt.contextErr,
tt.draining,
observe.SessionStartReasonMetadataStore,
)
if gotOutcome != tt.wantOutcome || gotReason != tt.wantReason {
t.Fatalf("operation result = (%q, %q), want (%q, %q)",
gotOutcome, gotReason, tt.wantOutcome, tt.wantReason)
}
})
}
if got := controlPlaneSessionStartOperationOutcome(operationErr, nil, false); got != "error" {
t.Fatalf("ordinary operation outcome = %q, want error", got)
}

func TestMissingOrgStackReason(t *testing.T) {
tests := []struct {
name string
state string
orgExists bool
want string
}{
{name: "missing org", orgExists: false, want: observe.SessionStartReasonClient},
{name: "legacy ready state", orgExists: true, state: "", want: observe.SessionStartReasonControlPlane},
{name: "ready warehouse", orgExists: true, state: string(configstore.ManagedWarehouseStateReady), want: observe.SessionStartReasonControlPlane},
{name: "failed provisioning", orgExists: true, state: string(configstore.ManagedWarehouseStateFailed), want: observe.SessionStartReasonLifecycle},
{name: "deleting", orgExists: true, state: string(configstore.ManagedWarehouseStateDeleting), want: observe.SessionStartReasonLifecycle},
{name: "resharding", orgExists: true, state: string(configstore.ManagedWarehouseStateResharding), want: observe.SessionStartReasonLifecycle},
{name: "pending", orgExists: true, state: string(configstore.ManagedWarehouseStatePending), want: observe.SessionStartReasonLifecycle},
}
if got := controlPlaneSessionStartOperationOutcome(operationErr, context.Canceled, true); got != "draining" {
t.Fatalf("drain-canceled operation outcome = %q, want draining", got)

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := missingOrgStackReason(tt.state, tt.orgExists); got != tt.want {
t.Fatalf("missingOrgStackReason(%q, %t) = %q, want %q", tt.state, tt.orgExists, got, tt.want)
}
})
}
}

Expand Down
Loading
Loading