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
8 changes: 6 additions & 2 deletions docs/invariants.md
Original file line number Diff line number Diff line change
Expand Up @@ -327,7 +327,9 @@ needing manual remediation aborts the pass rather than leaving storage half-conv
violated:* the first instance of a rolling deploy drops state the rest of the fleet is still
reading. *Enforced:* the per-dialect bootstrappers (`pkg/api/ensure_schema.go`,
`pkg/api/ensure_schema_postgres.go`), which the operator-facing storage schema surface calls rather
than reimplements (`pkg/api/storage_schema.go`).
than reimplements (`pkg/api/storage_schema.go`); the instance's own storage is the only target a
remote caller can address, and the deployment's permission is only ever widened, in the adapter that
answers for it (`pkg/serve/storage_schema.go`).

### AV-10: Anything the PR can do, the CLI can do

Expand Down Expand Up @@ -1370,7 +1372,9 @@ State storage must use an explicit connection and a different database name from
target in the same database family. This name check does not establish isolation for dynamically
resolved targets. Local hosting never permits destructive storage bootstrap.

*Enforced:* `pkg/serve/local.go`, `pkg/auth/local.go`, and `pkg/api/storage_isolation.go`.
*Enforced:* `pkg/serve/local.go`, `pkg/auth/local.go`, and `pkg/api/storage_isolation.go`; the request
that opts in to destructive storage statements is refused rather than honored on a locally hosted server
in `pkg/serve/storage_schema.go`.

### AZ-7: Local registration preserves existing work

Expand Down
6 changes: 5 additions & 1 deletion pkg/serve/local.go
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,11 @@ func RunLocal(ctx context.Context, config api.ServerConfig, local LocalOptions,
}
}()
opts = append(opts, func(o *options) { o.authorizer = authorizer })
srv, err := Build(ctx, &config, opts...)
// Local hosting travels with the server rather than only through this
// function's validation: the config key ValidateLocalConfig refuses is not
// the only way to ask for a destructive storage bootstrap now that an
// operator can ask for one per request (AZ-6).
srv, err := Build(ctx, &config, append(opts, withLocalHosting())...)
if err != nil {
return fmt.Errorf("build local runtime: %w", err)
}
Expand Down
25 changes: 25 additions & 0 deletions pkg/serve/local_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -104,3 +104,28 @@ func TestLocalListenerReservedBeforeBuild(t *testing.T) {
err = RunLocal(ctx, cfg, LocalOptions{Address: listener.Addr().String(), Token: strings.Repeat("a", 64)})
require.ErrorContains(t, err, "listen for local runtime")
}

// Every server the local runtime hosts is built as locally hosted, so the
// boundaries local hosting carries hold for the whole server rather than only
// for the config keys ValidateLocalConfig reads (AZ-6).
//
// The option is observed through the options value Build applies them to: the
// spy captures the pointer, and the local runtime's own option is applied to
// the same value afterwards. Build then fails on the unreadable TLS files,
// which needs no database.
func TestRunLocalBuildsALocallyHostedServer(t *testing.T) {
ctx, cancel := context.WithTimeout(t.Context(), time.Second)
defer cancel()
cfg := validLocalConfig()
missing := filepath.Join(t.TempDir(), "missing.pem")
cfg.PlanetScale.MTLS = &api.PlanetScaleMTLSConfig{CABundle: missing, ClientCert: missing, ClientKey: missing}
require.NoError(t, ValidateLocalConfig(&cfg))

var built *options
spy := func(o *options) { built = o }
err := RunLocal(ctx, cfg, LocalOptions{Address: "127.0.0.1:0", Token: strings.Repeat("a", 64)}, spy)

require.ErrorContains(t, err, "build local runtime")
require.NotNil(t, built, "Build applied the options it was handed")
assert.True(t, built.localHosted, "the local runtime claims local hosting for the server it builds")
}
186 changes: 156 additions & 30 deletions pkg/serve/serve.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,20 @@ type options struct {
date string
engines map[string]tern.EngineFactory
authorizer auth.Authorizer
// localHosted marks a server hosted by RunLocal. It is not a config field
// because it is not the operator's to set: it says which entry point built
// this server, and the local one carries boundaries the normal one does
// not (AZ-6).
localHosted bool
}

// withLocalHosting marks the server as hosted by the local runtime. It is
// unexported because the only caller that may claim it is RunLocal — a
// deployment that could assert local hosting through a config file or an
// embedder option could assert the boundaries AZ-6 grants it without accepting
// them, and a deployment that could deny it could drop them.
func withLocalHosting() Option {
return func(o *options) { o.localHosted = true }
}

// WithLogger sets the logger Run uses. A nil logger is ignored so Run keeps
Expand Down Expand Up @@ -111,6 +125,25 @@ func moduleVersion() string {
return versionFromBuildInfo(info)
}

// attributableVersion is the version a storage schema report may attribute
// this binary's embedded schema files to.
//
// The two uses of a version want opposite things from a build the module graph
// cannot name. A log field wants the sentinel present, so a query for the field
// finds the pod running an unidentifiable build instead of silently missing it.
// A report's attribution is prose an operator reads, and there "the schema
// embedded in unknown" reads as a release named unknown — a worse answer than
// the words the report already has for a build it cannot name, which say that
// these are the answering binary's own files without claiming whose build it
// is. So the sentinel stops here, at the one boundary where it would be read as
// a version rather than as a log value.
func attributableVersion(version string) string {
if version == unknownModuleVersion {
return ""
}
return version
}

// versionFromBuildInfo finds SchemaBot's version among the main module's
// dependencies. A replace directive wins, so a host pinning a fork or a local
// path is reported as what it actually runs rather than as the version it
Expand Down Expand Up @@ -303,6 +336,27 @@ type Server struct {
// probe never outlives the resolver and clients svc.Close tears down.
probeCancel context.CancelFunc
probeDone chan struct{}
// dialect is the storage database's family, resolved once at build time so
// every later storage operation — including the operator-facing storage
// schema surface — routes to the same family the bootstrap converged.
dialect schema.Dialect
// storageDSN is the DSN the storage pool was opened with. It names the
// database this instance booted against, which is the only database its
// storage schema surface may read or converge.
storageDSN string
// storageSchema is the one adapter that answers for that database, built
// once by Build and shared by the HTTP routes and the gRPC service, so the
// two surfaces cannot come to disagree about which storage they describe.
storageSchema tern.StorageSchemaService
// localHosted marks a server the local runtime hosts, which carries
// boundaries a normally hosted one does not (AZ-6).
localHosted bool
// version is the build's SchemaBot version as the logs carry it, which
// means it may be the unidentifiable-build sentinel. Storage schema
// reports attribute this binary's embedded files to it, and take it
// through attributableVersion on the way, because a report is prose and
// the sentinel is a log value.
version string
}

// registerPlanetScaleMTLS registers the configured planetscale.mtls
Expand Down Expand Up @@ -347,14 +401,17 @@ func Build(ctx context.Context, cfg *api.ServerConfig, opts ...Option) (*Server,
opt(&o)
}
logger := o.logger
if o.version == "" {
version := o.version
if version == "" {
// A host binary that embeds SchemaBot supplies its own logger and has no
// reason to know SchemaBot's version. Read it from the module graph, which
// is where an embedded dependency's version lives, so every log line
// identifies which SchemaBot the pod is running.
logger = logger.With("schemabot_version", moduleVersion())
// identifies which SchemaBot the pod is running — and so does every
// storage schema report, which attributes its embedded files to a build.
version = moduleVersion()
logger = logger.With("schemabot_version", version)
}
logger.Info("building server", "version", o.version, "commit", o.commit, "built", o.date)
logger.Info("building server", "version", version, "commit", o.commit, "built", o.date)

// Register PlanetScale mTLS before anything else so a worker with
// missing or unreadable certificate material fails startup immediately
Expand Down Expand Up @@ -396,7 +453,7 @@ func Build(ctx context.Context, cfg *api.ServerConfig, opts ...Option) (*Server,
// budget lets the pod wait the window out instead of crash-looping
// through it.
logger.Info("ensuring storage schema", "dialect", dialect)
db, err := bootStorage(ctx, cfg, dialect, logger)
db, storageDSN, err := bootStorage(ctx, cfg, dialect, logger)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -508,8 +565,7 @@ func Build(ctx context.Context, cfg *api.ServerConfig, opts ...Option) (*Server,
if router, ok := dataPlaneClient.(*tern.TargetRouter); ok {
targetResolver = router.Resolver()
}
success = true
return &Server{
srv := &Server{
cfg: cfg,
svc: svc,
storage: store,
Expand All @@ -520,7 +576,18 @@ func Build(ctx context.Context, cfg *api.ServerConfig, opts ...Option) (*Server,
telemetry: telemetry,
authz: authz,
engines: o.engines,
}, nil
dialect: dialect,
storageDSN: storageDSN,
version: version,
localHosted: o.localHosted,
}

if err := srv.registerStorageSchema(svc); err != nil {
return nil, err
}

success = true
return srv, nil
}

// Storage boot retry policy. The budget is sized so that even a final attempt
Expand All @@ -543,15 +610,18 @@ const inProcessWebhookDrainTimeout = 25 * time.Second
// failed attempts until the boot budget is spent. The DSN is re-resolved on
// every attempt so file-backed references pick up credentials rotated while
// the server waits.
func bootStorage(ctx context.Context, cfg *api.ServerConfig, dialect schema.Dialect, logger *slog.Logger) (*sql.DB, error) {
// It returns the DSN the successful attempt used alongside the pool, so the
// rest of the server can name the storage it actually booted against rather
// than re-resolving a value that may have moved since.
func bootStorage(ctx context.Context, cfg *api.ServerConfig, dialect schema.Dialect, logger *slog.Logger) (*sql.DB, string, error) {
deadline := time.Now().Add(storageBootRetryBudget)
for attempt := 1; ; attempt++ {
db, err := connectStorage(ctx, cfg, dialect, logger)
db, dsn, err := connectStorage(ctx, cfg, dialect, logger)
if err == nil {
return db, nil
return db, dsn, nil
}
if time.Until(deadline) < storageBootRetryInterval {
return nil, fmt.Errorf("storage not ready after %d attempts over %s: %w", attempt, storageBootRetryBudget, err)
return nil, "", fmt.Errorf("storage not ready after %d attempts over %s: %w", attempt, storageBootRetryBudget, err)
}
logger.Warn("storage not ready, retrying",
"attempt", attempt,
Expand All @@ -560,58 +630,106 @@ func bootStorage(ctx context.Context, cfg *api.ServerConfig, dialect schema.Dial
"error", err)
select {
case <-ctx.Done():
return nil, fmt.Errorf("storage boot canceled after %d attempts: %w", attempt, ctx.Err())
return nil, "", fmt.Errorf("storage boot canceled after %d attempts: %w", attempt, ctx.Err())
case <-time.After(storageBootRetryInterval):
}
}
}

// connectStorage runs a single storage boot attempt: resolve the DSN, apply
// the storage schema, open the pool, and verify it with a ping.
func connectStorage(ctx context.Context, cfg *api.ServerConfig, dialect schema.Dialect, logger *slog.Logger) (*sql.DB, error) {
// the storage schema, open the pool, and verify it with a ping. It returns the
// DSN it used so the caller holds the one this pool is dialing.
func connectStorage(ctx context.Context, cfg *api.ServerConfig, dialect schema.Dialect, logger *slog.Logger) (*sql.DB, string, error) {
const pingTimeout = 10 * time.Second
dsn, err := cfg.StorageDSN()
if err != nil {
return nil, fmt.Errorf("resolve storage DSN: %w", err)
return nil, "", fmt.Errorf("resolve storage DSN: %w", err)
}
if err := api.EnsureSchema(dsn, logger,
api.WithAllowDestructiveSchemaChanges(cfg.Storage.AllowDestructiveSchemaChanges),
api.WithPostgresStatementTimeout(cfg.Postgres.StatementTimeoutOrDefault()),
api.WithDialect(dialect)); err != nil {
return nil, fmt.Errorf("ensure storage schema: %w", err)
return nil, "", fmt.Errorf("ensure storage schema: %w", err)
}
db, err := openStoragePool(dialect, dsn, cfg)
db, err := openStoragePool(dialect, dsn, cfg, logger)
if err != nil {
return nil, fmt.Errorf("open storage database: %w", err)
return nil, "", fmt.Errorf("open storage database: %w", err)
}
pingCtx, cancel := context.WithTimeout(ctx, pingTimeout)
defer cancel()
if err := db.PingContext(pingCtx); err != nil {
utils.CloseAndLog(db)
return nil, fmt.Errorf("ping storage database: %w", err)
return nil, "", fmt.Errorf("ping storage database: %w", err)
}
return db, dsn, nil
Comment thread
aparajon marked this conversation as resolved.
}

// pinnedStorageDSN is the reload callback the storage pool re-resolves through
// after an authentication failure, narrowed to the one thing a reload is for.
//
// A rotated credential must be picked up without a restart, so the DSN is
// re-read from the live configuration. The database it names must not move,
// because everything downstream of the boot assumes it did not: the schema this
// server bootstrapped is on the database it booted against, and nothing
// bootstraps the new one. Without this guard an authentication failure is all
// it takes — the pool re-resolves, a config that now names another database
// answers the dial, and the server proceeds against storage whose schema it
// never converged.
//
// So the reload refuses a DSN whose address or database name has changed and
// keeps the pool on the database it booted against, failing the connection
// rather than silently relocating. Adopting new storage is a restart.
func pinnedStorageDSN(dialect schema.Dialect, bootDSN string, cfg *api.ServerConfig, logger *slog.Logger) func() (string, error) {
boot, bootErr := storageTargetFor(dialect, bootDSN)
return func() (string, error) {
if bootErr != nil {
return "", fmt.Errorf("read the storage target this server booted against: %w", bootErr)
}
next, err := cfg.StorageDSN()
if err != nil {
return "", fmt.Errorf("re-resolve storage DSN: %w", err)
}
resolved, err := storageTargetFor(dialect, next)
if err != nil {
return "", fmt.Errorf("read the storage target the current configuration names: %w", err)
}
if resolved != boot {
logger.Error("refusing to reconnect storage: the configured storage has moved since this server booted",
"booted_against", boot.String(),
"now_configured", resolved.String(),
"dialect", dialect)
return "", fmt.Errorf("the configured storage now names %s, but this server booted against %s; "+
"its storage schema was converged on the database it booted against, so restart it to adopt the new storage", resolved, boot)
}
return next, nil
}
return db, nil
}

// openStoragePool opens the long-lived reloadable storage pool for the
// configured dialect. Both connectors re-resolve the DSN through
// cfg.StorageDSN on authentication failure so a rotated storage credential
// is picked up without a restart. The dispatch fails closed: a dialect
// without a connector returns an error instead of dialing with another
// family's driver.
func openStoragePool(dialect schema.Dialect, dsn string, cfg *api.ServerConfig) (*sql.DB, error) {
// configured dialect. Both connectors re-resolve the DSN on authentication
// failure so a rotated storage credential is picked up without a restart. The
// dispatch fails closed: a dialect without a connector returns an error instead
// of dialing with another family's driver.
//
// The reload callback is built here rather than passed in. Every pool this
// opens must re-resolve through pinnedStorageDSN — a pool handed the raw
// resolver follows a rewritten secret to a database nothing bootstrapped — and
// a parameter is a place for the wrong callback to arrive. With none, there is
// no caller left to get it wrong.
func openStoragePool(dialect schema.Dialect, dsn string, cfg *api.ServerConfig, logger *slog.Logger) (*sql.DB, error) {
connectTimeout := cfg.Storage.Pool.ConnectTimeoutOrZero()
reload := pinnedStorageDSN(dialect, dsn, cfg, logger)
switch dialect {
case schema.DialectMySQL:
return mysqlconn.OpenReloadable(dsn, cfg.StorageDSN,
return mysqlconn.OpenReloadable(dsn, reload,
mysqlconn.WithConnectTimeout(connectTimeout))
case schema.DialectPostgres:
// The storage pool carries a statement budget of its own so steady-state
// storage queries run under a value SchemaBot states rather than
// whatever the platform imposed at the role or database level. It is
// the ordinary-query budget, not the bootstrap's DDL budget: this pool
// never executes DDL.
return postgresconn.OpenReloadable(dsn, cfg.StorageDSN,
return postgresconn.OpenReloadable(dsn, reload,
postgresconn.WithConnectTimeout(connectTimeout),
postgresconn.WithStatementTimeout(cfg.Postgres.StatementTimeoutOrDefault()))
default:
Expand Down Expand Up @@ -671,7 +789,15 @@ func (s *Server) RegisterGRPC(ctx context.Context, gs *grpc.Server) error {
s.svc.SetDefaultTernClient(built)
client = built
}
tern.NewServer(client, s.logger).Register(gs)
// The storage-schema service answers for this instance's own storage
// database, which is the only way a control plane can read it: a data
// plane's storage is reachable from the data plane, and the gRPC endpoint
// is the connection that already exists between the two.
opts := []tern.ServerOption{}
if s.storageSchema != nil {
opts = append(opts, tern.WithStorageSchemaService(s.storageSchema))
}
tern.NewServer(client, s.logger, opts...).Register(gs)
return nil
}

Expand Down
Loading
Loading