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
86 changes: 74 additions & 12 deletions internal/storage/db/cassandradb/audit_log.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package cassandradb
import (
"context"
"fmt"
"strings"
"time"

"github.com/gocql/gocql"
Expand Down Expand Up @@ -44,32 +45,82 @@ func (p *provider) ListAuditLogs(ctx context.Context, pagination *model.Paginati
queryBase := fmt.Sprintf("SELECT id, actor_id, actor_type, actor_email, action, resource_type, resource_id, ip_address, user_agent, metadata, created_at FROM %s", KeySpace+"."+schemas.Collections.AuditLog)
countBase := fmt.Sprintf("SELECT COUNT(*) FROM %s", KeySpace+"."+schemas.Collections.AuditLog)

whereClause := ""
// Every filter column below is backed by a secondary index (see
// provider.go), so equality restrictions need no ALLOW FILTERING — Scylla
// builds those indexes as materialized views and serves them directly.
//
// The created_at BOUNDS are different: a range on a non-primary-key column
// cannot be served by an index, so it forces ALLOW FILTERING and a scan.
// Only the timestamp bounds set that flag; an indexed-equality-only query
// keeps its existing index-served plan.
//
// ponytail: ALLOW FILTERING covers TWO cases, not just the rare one. The
// timestamp range is genuinely rare. Two or more equality filters is NOT —
// "this actor, this action" is an ordinary admin search, and it now scans
// where it previously errored outright. That is still the right trade here
// (admin-only endpoint, paginated, and an erroring filter combination is
// worse than a slow one), but it is a scan, not a free lunch. The upgrade
// path when it stops being cheap is a composite materialized view over the
// filter combinations that actually get used — not a bigger scan.
clauses := []string{}
filterValues := []interface{}{}
needsAllowFiltering := false

addEq := func(col string, v interface{}) {
clauses = append(clauses, col+"=?")
filterValues = append(filterValues, v)
}

if action, ok := filter["action"]; ok && action != "" {
whereClause += " WHERE action=?"
filterValues = append(filterValues, action)
addEq("action", action)
}
if actorID, ok := filter["actor_id"]; ok && actorID != "" {
if whereClause == "" {
whereClause += " WHERE actor_id=?"
} else {
whereClause += " AND actor_id=?"
}
filterValues = append(filterValues, actorID)
addEq("actor_id", actorID)
}
if resourceType, ok := filter["resource_type"]; ok && resourceType != "" {
addEq("resource_type", resourceType)
}
if resourceID, ok := filter["resource_id"]; ok && resourceID != "" {
addEq("resource_id", resourceID)
}
if fromTimestamp, ok := filter["from_timestamp"]; ok {
clauses = append(clauses, "created_at>=?")
filterValues = append(filterValues, fromTimestamp)
needsAllowFiltering = true
}
if toTimestamp, ok := filter["to_timestamp"]; ok {
clauses = append(clauses, "created_at<=?")
filterValues = append(filterValues, toTimestamp)
needsAllowFiltering = true
}

whereClause := ""
if len(clauses) > 0 {
whereClause = " WHERE " + strings.Join(clauses, " AND ")
}
// More than one restriction on non-primary-key columns cannot be served by a
// single index either, so Cassandra/Scylla requires the scan hint there too.
if len(clauses) > 1 {
needsAllowFiltering = true
}
if needsAllowFiltering {
whereClause += " ALLOW FILTERING"
}

// Count total — equality on indexed columns (action, actor_id) must not
// use ALLOW FILTERING; Scylla builds secondary indexes as materialized views.
countQuery := countBase + whereClause
err := p.db.Query(countQuery, filterValues...).Consistency(gocql.One).Scan(&paginationClone.Total)
if err != nil {
return nil, nil, err
}

// Fetch with pagination
query := queryBase + whereClause + fmt.Sprintf(" LIMIT %d", pagination.Limit+pagination.Offset)
// CQL grammar: LIMIT precedes ALLOW FILTERING, so the hint cannot simply be
// carried along on whereClause here.
query := queryBase + strings.TrimSuffix(whereClause, " ALLOW FILTERING") +
fmt.Sprintf(" LIMIT %d", pagination.Limit+pagination.Offset)
if needsAllowFiltering {
query += " ALLOW FILTERING"
}
scanner := p.db.Query(query, filterValues...).Iter().Scanner()
counter := int64(0)
for scanner.Next() {
Expand All @@ -87,6 +138,17 @@ func (p *provider) ListAuditLogs(ctx context.Context, pagination *model.Paginati
}
counter++
}
// A scan that dies part-way — read timeout, coordinator failure — ends
// Next() normally, so without this the call returns a TRUNCATED page with a
// nil error while paginationClone.Total (a separate query) reports the real
// count. An admin would see an incomplete audit trail with no signal it was
// cut short, which is the exact failure this table exists to prevent. Newly
// reachable here: the filters above can now produce ALLOW FILTERING scans,
// where a partial read is far likelier than on an index-served equality.
// DeleteAuditLogsBefore below already does this.
if err := scanner.Err(); err != nil {
return nil, nil, err
}

return auditLogs, &paginationClone, nil
}
Expand Down
73 changes: 51 additions & 22 deletions internal/storage/db/cassandradb/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -420,9 +420,20 @@ func NewProvider(cfg *config.Config, deps *Dependencies) (*provider, error) {
if err != nil {
return nil, err
}
// ScyllaDB builds secondary indexes asynchronously. Poll with a probe query
// that requires the actor_id index until it succeeds instead of a fixed sleep.
waitForCassandraIndexes(session, KeySpace, schemas.Collections.AuditLog, 30*time.Second)
auditLogResourceTypeIndex := fmt.Sprintf("CREATE INDEX IF NOT EXISTS authorizer_audit_log_resource_type ON %s.%s (resource_type)", KeySpace, schemas.Collections.AuditLog)
err = session.Query(auditLogResourceTypeIndex).Exec()
if err != nil {
return nil, err
}
auditLogResourceIDIndex := fmt.Sprintf("CREATE INDEX IF NOT EXISTS authorizer_audit_log_resource_id ON %s.%s (resource_id)", KeySpace, schemas.Collections.AuditLog)
err = session.Query(auditLogResourceIDIndex).Exec()
if err != nil {
return nil, err
}
// ScyllaDB builds secondary indexes asynchronously. Poll each indexed column
// until its index answers, instead of a fixed sleep.
waitForCassandraIndexes(session, KeySpace, schemas.Collections.AuditLog,
[]string{"actor_id", "action", "resource_type", "resource_id"}, 60*time.Second)

// Client table
clientCollectionQuery := fmt.Sprintf("CREATE TABLE IF NOT EXISTS %s.%s (id text, client_id text, kind text, name text, description text, client_secret text, allowed_scopes text, redirect_uris text, grant_types text, token_endpoint_auth_method text, is_active boolean, org_id text, created_at bigint, updated_at bigint, PRIMARY KEY (id))", KeySpace, schemas.Collections.Client)
Expand Down Expand Up @@ -634,7 +645,13 @@ func NewProvider(cfg *config.Config, deps *Dependencies) (*provider, error) {
// fail until the index is ready.
func waitForCassandraSecondaryIndex(session *cansandraDriver.Session, keyspace, table, column string, timeout time.Duration) {
probe := fmt.Sprintf("SELECT id FROM %s.%s WHERE %s='' LIMIT 1 ALLOW FILTERING", keyspace, table, column)
deadline := time.Now().Add(timeout)
pollUntilQuerySucceeds(session, probe, time.Now().Add(timeout))
}

// pollUntilQuerySucceeds retries probe with a capped linear backoff until it
// succeeds or deadline passes. Shared so the backoff cannot drift between the
// single-column and multi-column waits.
func pollUntilQuerySucceeds(session *cansandraDriver.Session, probe string, deadline time.Time) {
delay := 500 * time.Millisecond
for {
if err := session.Query(probe).Exec(); err == nil {
Expand All @@ -650,24 +667,36 @@ func waitForCassandraSecondaryIndex(session *cansandraDriver.Session, keyspace,
}
}

// waitForCassandraIndexes polls a probe query that requires the actor_id secondary
// index until it succeeds or the timeout is reached. ScyllaDB builds secondary
// indexes asynchronously; queries on indexed columns fail until the index is ready.
func waitForCassandraIndexes(session *cansandraDriver.Session, keyspace, table string, timeout time.Duration) {
probe := fmt.Sprintf("SELECT id FROM %s.%s WHERE actor_id='' LIMIT 1", keyspace, table)
deadline := time.Now().Add(timeout)
delay := 500 * time.Millisecond
for {
if err := session.Query(probe).Exec(); err == nil {
return
}
if time.Now().After(deadline) {
return
}
time.Sleep(delay)
if delay < 3*time.Second {
delay += 500 * time.Millisecond
}
// waitForCassandraIndexes probes each indexed column until its secondary index
// answers, or until that column's share of the timeout runs out. ScyllaDB builds
// secondary indexes asynchronously and queries on an indexed column fail until
// the index is ready.
//
// Every column is probed, not just one: Scylla builds indexes concurrently, so
// the last one created is not necessarily the last one ready, and waiting on a
// single column would let a query against another index run before that index
// exists.
//
// The budget is split PER COLUMN rather than shared. With one shared deadline a
// slow first index consumes the whole budget and the remaining columns are never
// probed at all — turning a startup race on one column into a startup race on
// every other, which is worse than the single-column wait this replaced. An
// exhausted column moves on rather than abandoning the rest.
//
// Best-effort by design: this is startup smoothing, not a correctness gate, so an
// expired probe is not an error.
func waitForCassandraIndexes(session *cansandraDriver.Session, keyspace, table string, columns []string, timeout time.Duration) {
if len(columns) == 0 {
return
}
perColumn := timeout / time.Duration(len(columns))
for _, col := range columns {
// No ALLOW FILTERING here, deliberately: the hint would let the probe
// succeed whether or not the index exists, so the wait would return
// immediately and detect nothing. The probe has to be a query that
// FAILS until the index is ready.
probe := fmt.Sprintf("SELECT id FROM %s.%s WHERE %s='' LIMIT 1", keyspace, table, col)
pollUntilQuerySucceeds(session, probe, time.Now().Add(perColumn))
}
}

Expand Down
36 changes: 29 additions & 7 deletions internal/storage/db/couchbase/audit_log.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package couchbase
import (
"context"
"fmt"
"strings"
"time"

"github.com/couchbase/gocb/v2"
Expand Down Expand Up @@ -43,19 +44,40 @@ func (p *provider) ListAuditLogs(ctx context.Context, pagination *model.Paginati
params["offset"] = paginationClone.Offset
params["limit"] = paginationClone.Limit

whereClause := ""
// Every filter the ListAuditLogRequest API accepts is applied here. Leaving
// any of them out silently returns unfiltered rows: the caller cannot tell a
// filter was dropped, so an auditor narrowing to one resource or one time
// window would read the result as authoritative.
clauses := []string{}
if action, ok := filter["action"]; ok && action != "" {
whereClause += " WHERE action=$action"
clauses = append(clauses, "action=$action")
params["action"] = action
}
if actorID, ok := filter["actor_id"]; ok && actorID != "" {
if whereClause == "" {
whereClause += " WHERE actor_id=$actorID"
} else {
whereClause += " AND actor_id=$actorID"
}
clauses = append(clauses, "actor_id=$actorID")
params["actorID"] = actorID
}
if resourceType, ok := filter["resource_type"]; ok && resourceType != "" {
clauses = append(clauses, "resource_type=$resourceType")
params["resourceType"] = resourceType
}
if resourceID, ok := filter["resource_id"]; ok && resourceID != "" {
clauses = append(clauses, "resource_id=$resourceID")
params["resourceID"] = resourceID
}
if fromTimestamp, ok := filter["from_timestamp"]; ok {
clauses = append(clauses, "created_at>=$fromTimestamp")
params["fromTimestamp"] = fromTimestamp
}
if toTimestamp, ok := filter["to_timestamp"]; ok {
clauses = append(clauses, "created_at<=$toTimestamp")
params["toTimestamp"] = toTimestamp
}

whereClause := ""
if len(clauses) > 0 {
whereClause = " WHERE " + strings.Join(clauses, " AND ")
}

// Count with filters applied
countQuery := fmt.Sprintf("SELECT COUNT(*) as count FROM %s.%s%s",
Expand Down
42 changes: 40 additions & 2 deletions internal/storage/db/couchbase/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,17 @@ func isTransientQueryErr(err error) bool {
if err == nil {
return false
}
// errors.Is on the shared parent, not a substring match on the message.
// gocbcore defines ErrAmbiguousTimeout as "ambiguous timeout" and
// ErrUnambiguousTimeout as "unambiguous timeout", both wrapping ErrTimeout.
// A substring test for "unambiguous timeout" therefore MISSES the ambiguous
// one — and a CREATE INDEX that outruns the client is precisely the
// ambiguous case, because the server may well have accepted it. Missing it
// meant no retry, and a non-retried index error aborts provider
// construction, i.e. the server does not start.
if errors.Is(err, gocb.ErrTimeout) {
return true
}
msg := strings.ToLower(err.Error())
return strings.Contains(msg, "eof") ||
strings.Contains(msg, "connection reset") ||
Expand All @@ -155,11 +166,31 @@ func isTransientQueryErr(err error) bool {
}

// execIndexQuery runs a CREATE INDEX statement with retries for transient query-service errors.
//
// The per-attempt timeout is bounded to a third of the overall budget, and is
// passed explicitly rather than left to gocb's default. gocb's default
// QueryTimeout is 75s while this budget defaults to 30s, so a single attempt
// could outrun the whole deadline and return past it — the retry that would
// then see "already exists" never ran, and the resulting error aborts provider
// construction (the caller does `return nil, err`), so the server fails to
// boot.
//
// That is not hypothetical on upgrade: Couchbase CREATE INDEX without
// defer_build is synchronous, and a newly added index on a collection that
// already holds months of audit rows builds over all of them. Timing out the
// attempt is harmless — the index definition is registered server-side and the
// build continues — so the next attempt returns "already exists" and startup
// proceeds while the backfill finishes in the background. That matches the
// async DDL behaviour the Cassandra provider already relies on.
func execIndexQuery(scope *gocb.Scope, query string, timeout time.Duration) error {
deadline := time.Now().Add(timeout)
perAttempt := timeout / 3
if perAttempt <= 0 {
perAttempt = timeout
}
delay := 500 * time.Millisecond
for {
_, err := scope.Query(query, nil)
_, err := scope.Query(query, &gocb.QueryOptions{Timeout: perAttempt})
if err == nil {
return nil
}
Expand Down Expand Up @@ -287,7 +318,14 @@ func getIndex(scopeName string) map[string][]string {
auditLogIndex1 := fmt.Sprintf("CREATE INDEX AuditLogActorIdIndex ON %s.%s(actor_id)", scopeName, schemas.Collections.AuditLog)
auditLogIndex2 := fmt.Sprintf("CREATE INDEX AuditLogActionIndex ON %s.%s(action)", scopeName, schemas.Collections.AuditLog)
auditLogIndex3 := fmt.Sprintf("CREATE INDEX AuditLogCreatedAtIndex ON %s.%s(created_at)", scopeName, schemas.Collections.AuditLog)
indices[schemas.Collections.AuditLog] = []string{auditLogIndex1, auditLogIndex2, auditLogIndex3}
// resource_type / resource_id are filterable via ListAuditLogs. Without these
// the predicates still return correct rows, but N1QL falls back to the
// CREATE PRIMARY INDEX scan of the fastest-growing collection in the scope,
// while Cassandra serves the same filter from an index — the "O(1) in one
// backend, full scan in another" parity bug AGENTS.md calls out.
auditLogIndex4 := fmt.Sprintf("CREATE INDEX AuditLogResourceTypeIndex ON %s.%s(resource_type)", scopeName, schemas.Collections.AuditLog)
auditLogIndex5 := fmt.Sprintf("CREATE INDEX AuditLogResourceIdIndex ON %s.%s(resource_id)", scopeName, schemas.Collections.AuditLog)
indices[schemas.Collections.AuditLog] = []string{auditLogIndex1, auditLogIndex2, auditLogIndex3, auditLogIndex4, auditLogIndex5}

// TrustedIssuer indexes
trustedIssuerIndex1 := fmt.Sprintf("CREATE INDEX TrustedIssuerIssuerURLIndex ON %s.%s(issuer_url)", scopeName, schemas.Collections.TrustedIssuer)
Expand Down
56 changes: 56 additions & 0 deletions internal/storage/db/couchbase/provider_index_retry_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
package couchbase

import (
"errors"
"fmt"
"testing"

"github.com/couchbase/gocb/v2"
"github.com/stretchr/testify/assert"
)

// A CREATE INDEX that outruns the client is the AMBIGUOUS timeout case — the
// server may have accepted it. If that is not treated as transient, the retry
// that would see "already exists" never runs, and a non-retried index error
// aborts provider construction: the server does not boot.
//
// The original check substring-matched "unambiguous timeout", which does not
// match gocbcore's "ambiguous timeout". This pins both.
func TestIsTransientQueryErr_CoversBothTimeoutKinds(t *testing.T) {
for _, tc := range []struct {
name string
err error
}{
{"ambiguous", gocb.ErrAmbiguousTimeout},
{"unambiguous", gocb.ErrUnambiguousTimeout},
{"wrapped ambiguous", fmt.Errorf("create index: %w", gocb.ErrAmbiguousTimeout)},
} {
t.Run(tc.name, func(t *testing.T) {
assert.True(t, isTransientQueryErr(tc.err),
"%v must be retryable, else index creation aborts startup", tc.err)
})
}
}

// Guard the premise of the fix: the two sentinels really do share ErrTimeout,
// and their messages really do differ in a way substring matching gets wrong.
func TestGocbTimeoutSentinels_ShareParentButNotMessage(t *testing.T) {
assert.True(t, errors.Is(gocb.ErrAmbiguousTimeout, gocb.ErrTimeout))
assert.True(t, errors.Is(gocb.ErrUnambiguousTimeout, gocb.ErrTimeout))
assert.NotContains(t, gocb.ErrAmbiguousTimeout.Error(), "unambiguous timeout",
"if this ever contains the longer string, the old substring check was fine")
}

func TestIsTransientQueryErr_NonTransientStaysFatal(t *testing.T) {
assert.False(t, isTransientQueryErr(nil))
assert.False(t, isTransientQueryErr(errors.New("syntax error in CREATE INDEX")))
}

// "already exists" short-circuits before the transient check, so a re-run
// against a database that already has the index is a no-op rather than a retry
// loop. This is what makes repeated boots safe.
func TestIsIndexExistsErr_TolerantOnRestart(t *testing.T) {
assert.True(t, isIndexExistsErr("Index AuditLogResourceIdIndex already exists"))
assert.True(t, isIndexExistsErr("The index #primary already exists"))
assert.False(t, isIndexExistsErr("ambiguous timeout"))
}
Loading
Loading