mirror of
https://github.com/pocket-id/pocket-id.git
synced 2026-09-30 23:09:05 +02:00
feat: add FRANCIS_HOST to connect to a standalone Francis runtime
FRANCIS_HOST decides where the Francis actor runtime lives. When it is empty or set to "embedded" (the default), Pocket ID starts the runtime inside its own process, backed by its own database. Any other value is the address, or a comma-separated list of addresses, of a standalone Francis runtime. Pocket ID then connects to it as a remote actor host and starts no embedded runtime. Note: connecting to a standalone runtime also needs FRANCIS_HOST_PSK or FRANCIS_HOST_JWT, and optionally (but recommended) FRANCIS_CA. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01DMgoTZtznSjP4SHTaHbRen
This commit is contained in:
co-authored by
Claude Opus 5
parent
81cb290bed
commit
cea1267f37
@@ -13,7 +13,9 @@ import (
|
||||
"github.com/italypaleale/francis/components"
|
||||
"github.com/italypaleale/francis/components/postgres"
|
||||
"github.com/italypaleale/francis/components/sqlite"
|
||||
francishost "github.com/italypaleale/francis/host"
|
||||
"github.com/italypaleale/francis/host/local"
|
||||
"github.com/italypaleale/francis/host/remote"
|
||||
"github.com/jackc/pgx/v5/pgxpool"
|
||||
"gorm.io/gorm"
|
||||
|
||||
@@ -24,6 +26,14 @@ import (
|
||||
"github.com/pocket-id/pocket-id/backend/internal/utils/crypto"
|
||||
)
|
||||
|
||||
// ErrRemoteFrancisRuntime is returned by the helpers that reach the actor data through Pocket ID's own database when FRANCIS_HOST points to a standalone Francis runtime
|
||||
// That runtime owns the actor data instead, so it can only be reached through the runtime itself
|
||||
var ErrRemoteFrancisRuntime = errors.New("the actor data is owned by the standalone Francis runtime configured in FRANCIS_HOST, and is not stored in Pocket ID's database")
|
||||
|
||||
// ErrEmbeddedFrancisRuntime is returned by the helpers that reach the actor data through a standalone Francis runtime when Pocket ID runs an embedded one
|
||||
// There is no runtime to connect to in that case, and the actor data is in Pocket ID's own database
|
||||
var ErrEmbeddedFrancisRuntime = errors.New("the actor runtime is embedded in Pocket ID, so there is no standalone Francis runtime to connect to")
|
||||
|
||||
type NewActorsOpts struct {
|
||||
Postgres *pgxpool.Pool
|
||||
|
||||
@@ -34,57 +44,25 @@ type NewActorsOpts struct {
|
||||
FileStorage storage.FileStorage
|
||||
}
|
||||
|
||||
func NewActors(o NewActorsOpts) (*local.Host, map[string]*ratelimit.RateLimitService, error) {
|
||||
log := slog.Default()
|
||||
func NewActors(o NewActorsOpts) (francishost.Host, map[string]*ratelimit.RateLimitService, error) {
|
||||
log := slog.Default().With("scope", "actor-host")
|
||||
|
||||
// Derive a PSK from the global encryption key
|
||||
// The runtime PSK derives the cluster CA used for host-to-host mTLS
|
||||
psk, err := o.getPSK()
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("failed to derive PSK: %w", err)
|
||||
// Create the actor host for the configured topology
|
||||
// The embedded runtime keeps the actor data in Pocket ID's own database, while a standalone Francis runtime owns it instead and coordinates every host that connects to it
|
||||
var (
|
||||
h francishost.Host
|
||||
err error
|
||||
)
|
||||
if o.EnvConfig.HasEmbeddedFrancisRuntime() {
|
||||
log.Debug("Starting the embedded Francis runtime")
|
||||
h, err = o.newEmbeddedHost(log)
|
||||
} else {
|
||||
log.Info("Connecting to a standalone Francis runtime", slog.Any("addresses", o.EnvConfig.FrancisAddresses))
|
||||
h, err = o.newRemoteHost(log)
|
||||
}
|
||||
|
||||
// Derive the cluster host limit from the HA setting
|
||||
// With HA disabled the cluster is capped at a single replica
|
||||
maxHosts := 1
|
||||
if o.EnvConfig.HAEnabled {
|
||||
// 0 = no cap
|
||||
maxHosts = 0
|
||||
}
|
||||
|
||||
// Options for the host
|
||||
opts := []local.HostOption{
|
||||
local.WithAddress(net.JoinHostPort(o.EnvConfig.ActorsHost, o.EnvConfig.ActorsPort)),
|
||||
local.WithLogger(log.With("scope", "actor-host")),
|
||||
local.WithRuntimePSKs(psk),
|
||||
local.WithShutdownGracePeriod(10 * time.Second),
|
||||
local.WithMaxHosts(maxHosts),
|
||||
local.WithHostHealthCheckDeadline(ActorsHostHealthCheckDeadline(o.EnvConfig.HAEnabled)),
|
||||
}
|
||||
|
||||
// With a single active host the relaxed alarm intervals reduce database load
|
||||
// The longer lease duration also means fewer lease renewals, since Francis renews a lease 10s before it expires (no other host can claim the alarm anyways)
|
||||
// When HA is enabled these are dropped so Francis uses its tighter defaults, which distribute alarm work and fail over faster across multiple hosts
|
||||
if !o.EnvConfig.HAEnabled {
|
||||
opts = append(opts,
|
||||
local.WithAlarmsPollInterval(5*time.Minute),
|
||||
local.WithAlarmsFetchAheadInterval(5*time.Minute),
|
||||
local.WithAlarmsLeaseDuration(180*time.Second),
|
||||
)
|
||||
}
|
||||
|
||||
// Add the database connection
|
||||
providerOpt, err := o.getProviderOption()
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
opts = append(opts, providerOpt)
|
||||
|
||||
// Create a new actor host
|
||||
h, err := local.NewHost(opts...)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("failed to create actor host: %w", err)
|
||||
}
|
||||
|
||||
// Add all cron jobs
|
||||
err = o.registerCronJobs(h)
|
||||
@@ -107,6 +85,108 @@ func NewActors(o NewActorsOpts) (*local.Host, map[string]*ratelimit.RateLimitSer
|
||||
return h, rateLimitServices, nil
|
||||
}
|
||||
|
||||
// newEmbeddedHost creates the actor host that runs the Francis runtime inside the Pocket ID process, backed by Pocket ID's own database
|
||||
func (o *NewActorsOpts) newEmbeddedHost(log *slog.Logger) (*local.Host, error) {
|
||||
// Derive a PSK from the global encryption key
|
||||
// The runtime PSK derives the cluster CA used for host-to-host mTLS
|
||||
psk, err := o.getPSK()
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to derive PSK: %w", err)
|
||||
}
|
||||
|
||||
// Derive the cluster host limit from the HA setting
|
||||
// With HA disabled the cluster is capped at a single replica
|
||||
maxHosts := 1
|
||||
if o.EnvConfig.HAEnabled {
|
||||
// 0 = no cap
|
||||
maxHosts = 0
|
||||
}
|
||||
|
||||
// Options for the host
|
||||
opts := []local.HostOption{
|
||||
local.WithAddress(net.JoinHostPort(o.EnvConfig.ActorsHost, o.EnvConfig.ActorsPort)),
|
||||
local.WithLogger(log),
|
||||
local.WithRuntimePSKs(psk),
|
||||
local.WithShutdownGracePeriod(10 * time.Second),
|
||||
local.WithMaxHosts(maxHosts),
|
||||
local.WithHostHealthCheckDeadline(ActorsHostHealthCheckDeadline(o.EnvConfig.HAEnabled)),
|
||||
}
|
||||
|
||||
// With a single active host the relaxed alarm intervals reduce database load
|
||||
// The longer lease duration also means fewer lease renewals, since Francis renews a lease 10s before it expires (no other host can claim the alarm anyways)
|
||||
// When HA is enabled these are dropped so Francis uses its tighter defaults, which distribute alarm work and fail over faster across multiple hosts
|
||||
if !o.EnvConfig.HAEnabled {
|
||||
opts = append(opts,
|
||||
local.WithAlarmsPollInterval(5*time.Minute),
|
||||
local.WithAlarmsFetchAheadInterval(5*time.Minute),
|
||||
local.WithAlarmsLeaseDuration(180*time.Second),
|
||||
)
|
||||
}
|
||||
|
||||
// Add the database connection
|
||||
providerOpt, err := o.getProviderOption()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
opts = append(opts, providerOpt)
|
||||
|
||||
h, err := local.NewHost(opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create actor host: %w", err)
|
||||
}
|
||||
|
||||
return h, nil
|
||||
}
|
||||
|
||||
// newRemoteHost creates the actor host that connects to a standalone Francis runtime
|
||||
// The runtime owns the actor state, placement, and alarms, so none of the embedded runtime's database and clustering options apply here
|
||||
// That includes the cap on the number of hosts in the cluster, which the runtime enforces through its own "maxHosts" setting: Pocket ID cannot limit itself to a single replica from this side
|
||||
func (o *NewActorsOpts) newRemoteHost(log *slog.Logger) (*remote.Host, error) {
|
||||
opts := append(
|
||||
remoteConnectionOptions(o.EnvConfig, log),
|
||||
// Actors placed on this host are invoked by its peers at this address, which is also the one it advertises to the runtime
|
||||
remote.WithAddress(net.JoinHostPort(o.EnvConfig.ActorsHost, o.EnvConfig.ActorsPort)),
|
||||
remote.WithShutdownGracePeriod(10*time.Second),
|
||||
)
|
||||
|
||||
h, err := remote.NewHost(opts...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create remote actor host: %w", err)
|
||||
}
|
||||
|
||||
return h, nil
|
||||
}
|
||||
|
||||
// remoteConnectionOptions builds the options that address and authenticate Pocket ID to a standalone Francis runtime
|
||||
// Both the actor host and the short-lived client the CLI commands use go through here, so they always present the same identity to the same cluster
|
||||
func remoteConnectionOptions(envConfig *common.EnvConfigSchema, log *slog.Logger) []remote.HostOption {
|
||||
opts := []remote.HostOption{
|
||||
remote.WithLogger(log),
|
||||
remote.WithRuntimeAddresses(envConfig.FrancisAddresses...),
|
||||
}
|
||||
|
||||
// The configuration is validated to carry exactly one bootstrap method, so the first match is the one the operator chose
|
||||
switch {
|
||||
case len(envConfig.FrancisHostPSK) > 0:
|
||||
opts = append(opts, remote.WithHostBootstrapPSK(envConfig.FrancisHostPSK))
|
||||
case envConfig.FrancisHostJWTFile != "":
|
||||
// Francis re-reads the file on every connection, so a rotated token is picked up without restarting Pocket ID
|
||||
opts = append(opts, remote.WithHostBootstrapJWTFile(envConfig.FrancisHostJWTFile))
|
||||
case envConfig.FrancisHostJWT != "":
|
||||
opts = append(opts, remote.WithHostBootstrapJWT(envConfig.FrancisHostJWT))
|
||||
}
|
||||
|
||||
// Pinning the cluster CA lets Pocket ID verify the runtime on its very first connection
|
||||
// Francis requires the trust decision to be explicit, so without a pinned CA we have to opt into trusting the certificate served on first use, which it warns about
|
||||
if len(envConfig.FrancisCA) > 0 {
|
||||
opts = append(opts, remote.WithPinnedCA(envConfig.FrancisCA))
|
||||
} else {
|
||||
opts = append(opts, remote.WithUnsafeNoPinnedCA())
|
||||
}
|
||||
|
||||
return opts
|
||||
}
|
||||
|
||||
// Derive a PSK from the global encryption key
|
||||
func (o *NewActorsOpts) getPSK() ([]byte, error) {
|
||||
// This is tied to the instance ID of the Pocket ID deployment/cluster
|
||||
@@ -117,7 +197,12 @@ func (o *NewActorsOpts) getPSK() ([]byte, error) {
|
||||
// NewActorStateStore creates a minimal actor host that can read and write actor state directly, without joining the cluster or binding a network port.
|
||||
// It's meant for short-lived contexts such as CLI commands that need to persist actor state (for example, one-time access tokens) without running the full actor host.
|
||||
// The returned host must NOT be Run(): only direct state operations (Get/Set/Delete on state) are supported, and they require the actor state tables to already exist, which is the case whenever the server has run at least once against this database.
|
||||
// It only works with the embedded runtime, since the actor state then lives in Pocket ID's own database: with a standalone Francis runtime it returns ErrRemoteFrancisRuntime.
|
||||
func NewActorStateStore(o NewActorsOpts) (*local.Host, error) {
|
||||
if !o.EnvConfig.HasEmbeddedFrancisRuntime() {
|
||||
return nil, ErrRemoteFrancisRuntime
|
||||
}
|
||||
|
||||
providerOpt, err := o.getProviderOption()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
@@ -234,7 +319,7 @@ func (o *NewActorsOpts) getProviderOption() (local.HostOption, error) {
|
||||
}
|
||||
}
|
||||
|
||||
func (o *NewActorsOpts) registerCronJobs(host *local.Host) (err error) {
|
||||
func (o *NewActorsOpts) registerCronJobs(host francishost.Host) (err error) {
|
||||
// In test mode, we do not register anything
|
||||
if common.EnvConfig.AppEnv == "test" {
|
||||
return nil
|
||||
@@ -271,7 +356,7 @@ func (o *NewActorsOpts) registerCronJobs(host *local.Host) (err error) {
|
||||
|
||||
// registerRateLimiters creates a built-in rate-limit actor for each middleware policy and returns both the created actors (keyed by policy name) and the host options to register them
|
||||
// Unlike cron jobs, rate limiters keep no durable state, so they are registered in every environment
|
||||
func (o *NewActorsOpts) registerRateLimiters(host *local.Host) (actors map[string]*ratelimit.RateLimit, err error) {
|
||||
func (o *NewActorsOpts) registerRateLimiters(host francishost.Host) (actors map[string]*ratelimit.RateLimit, err error) {
|
||||
policies := middleware.RateLimitPolicies()
|
||||
actors = make(map[string]*ratelimit.RateLimit, len(policies))
|
||||
for _, p := range policies {
|
||||
|
||||
@@ -2,10 +2,23 @@ package bootstrap
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/ed25519"
|
||||
"crypto/rand"
|
||||
"crypto/x509"
|
||||
"crypto/x509/pkix"
|
||||
"encoding/hex"
|
||||
"encoding/pem"
|
||||
"log/slog"
|
||||
"math/big"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
francishost "github.com/italypaleale/francis/host"
|
||||
"github.com/italypaleale/francis/host/local"
|
||||
"github.com/italypaleale/francis/host/remote"
|
||||
"github.com/libtnb/sqlite"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/gorm"
|
||||
@@ -31,6 +44,151 @@ func TestNewActorsOptsGetPSKUsesStableValue(t *testing.T) {
|
||||
require.Equalf(t, expected, actual, "actual result: %s", actual)
|
||||
}
|
||||
|
||||
// TestNewActorsSelectsTopology covers the branch that FRANCIS_HOST drives: with no standalone runtime configured Pocket ID starts an embedded one, and otherwise it connects to the addresses it was given.
|
||||
func TestNewActorsSelectsTopology(t *testing.T) {
|
||||
// The actor host is created but never run, so the database only has to exist
|
||||
newDB := func(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
|
||||
dbPath := filepath.Join(t.TempDir(), "pocket-id.db")
|
||||
db, err := gorm.Open(sqlite.Open("file:"+dbPath+"?_pragma=foreign_keys(1)"), &gorm.Config{})
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
sqlDB, dbErr := db.DB()
|
||||
if dbErr == nil {
|
||||
_ = sqlDB.Close()
|
||||
}
|
||||
})
|
||||
|
||||
return db
|
||||
}
|
||||
|
||||
// registerCronJobs reads the app environment off the global config and skips every job in test mode, which keeps each case down to the host itself and the rate limiters
|
||||
baseConfig := func(t *testing.T) *common.EnvConfigSchema {
|
||||
t.Helper()
|
||||
|
||||
originalAppEnv := common.EnvConfig.AppEnv
|
||||
common.EnvConfig.AppEnv = common.AppEnvTest
|
||||
t.Cleanup(func() {
|
||||
common.EnvConfig.AppEnv = originalAppEnv
|
||||
})
|
||||
|
||||
return &common.EnvConfigSchema{
|
||||
AppEnv: common.AppEnvTest,
|
||||
EncryptionKey: []byte("test-encryption-key"),
|
||||
ActorsHost: "127.0.0.1",
|
||||
ActorsPort: "1414",
|
||||
}
|
||||
}
|
||||
|
||||
t.Run("embedded runtime by default", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
|
||||
h, rateLimitServices, err := NewActors(NewActorsOpts{
|
||||
EnvConfig: cfg,
|
||||
InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876",
|
||||
DB: newDB(t),
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.IsType(t, &local.Host{}, h)
|
||||
require.NotEmpty(t, rateLimitServices)
|
||||
})
|
||||
|
||||
t.Run("remote runtime when addresses are configured", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.FrancisAddresses = []string{"runtime-1.example.com:8443", "runtime-2.example.com:8443"}
|
||||
cfg.FrancisHostPSK = []byte("bootstrap-psk-that-is-long-enough")
|
||||
|
||||
// No database is passed, since a standalone runtime owns the actor data and the remote host must not reach for Pocket ID's own database
|
||||
h, rateLimitServices, err := NewActors(NewActorsOpts{
|
||||
EnvConfig: cfg,
|
||||
InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876",
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.IsType(t, &remote.Host{}, h)
|
||||
require.NotEmpty(t, rateLimitServices)
|
||||
})
|
||||
|
||||
// Each bootstrap method has to produce a host Francis accepts, which is the only part of the remote wiring that can be checked without a runtime to connect to
|
||||
t.Run("every bootstrap method builds a valid remote host", func(t *testing.T) {
|
||||
jwtFile := filepath.Join(t.TempDir(), "token")
|
||||
require.NoError(t, os.WriteFile(jwtFile, []byte("header.payload.signature"), 0600))
|
||||
|
||||
for name, apply := range map[string]func(cfg *common.EnvConfigSchema){
|
||||
"PSK": func(cfg *common.EnvConfigSchema) { cfg.FrancisHostPSK = []byte("bootstrap-psk-that-is-long-enough") },
|
||||
"JWT": func(cfg *common.EnvConfigSchema) { cfg.FrancisHostJWT = "header.payload.signature" },
|
||||
"JWT file": func(cfg *common.EnvConfigSchema) { cfg.FrancisHostJWTFile = jwtFile },
|
||||
} {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.FrancisAddresses = []string{"runtime-1.example.com:8443"}
|
||||
apply(cfg)
|
||||
|
||||
opts := NewActorsOpts{EnvConfig: cfg, InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876"}
|
||||
h, err := opts.newRemoteHost(slog.New(slog.DiscardHandler))
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, h)
|
||||
})
|
||||
}
|
||||
})
|
||||
|
||||
// The E2E suite runs the remote variant without a pinned CA, so this is the only place the pinning branch is exercised
|
||||
t.Run("pinning the cluster CA builds a valid remote host", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.FrancisAddresses = []string{"runtime-1.example.com:8443"}
|
||||
cfg.FrancisHostPSK = []byte("bootstrap-psk-that-is-long-enough")
|
||||
cfg.FrancisCA = testCAPEM(t)
|
||||
|
||||
opts := NewActorsOpts{EnvConfig: cfg, InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876"}
|
||||
h, err := opts.newRemoteHost(slog.New(slog.DiscardHandler))
|
||||
require.NoError(t, err)
|
||||
require.NotNil(t, h)
|
||||
})
|
||||
|
||||
t.Run("an unparsable cluster CA is rejected", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.FrancisAddresses = []string{"runtime-1.example.com:8443"}
|
||||
cfg.FrancisHostPSK = []byte("bootstrap-psk-that-is-long-enough")
|
||||
cfg.FrancisCA = []byte("-----BEGIN CERTIFICATE-----\nnot a certificate\n-----END CERTIFICATE-----")
|
||||
|
||||
opts := NewActorsOpts{EnvConfig: cfg, InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876"}
|
||||
_, err := opts.newRemoteHost(slog.New(slog.DiscardHandler))
|
||||
require.Error(t, err)
|
||||
})
|
||||
|
||||
t.Run("no bootstrap method is rejected by Francis", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.FrancisAddresses = []string{"runtime-1.example.com:8443"}
|
||||
|
||||
opts := NewActorsOpts{EnvConfig: cfg, InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876"}
|
||||
_, err := opts.newRemoteHost(slog.New(slog.DiscardHandler))
|
||||
require.Error(t, err)
|
||||
})
|
||||
|
||||
t.Run("the actor client requires a standalone runtime", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
|
||||
err := WithActorClient(t.Context(), cfg, func(context.Context, francishost.Host) error {
|
||||
t.Fatal("the callback must not run without a standalone runtime")
|
||||
return nil
|
||||
})
|
||||
require.ErrorIs(t, err, ErrEmbeddedFrancisRuntime)
|
||||
})
|
||||
|
||||
t.Run("state store is unavailable with a remote runtime", func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.FrancisAddresses = []string{"runtime-1.example.com:8443"}
|
||||
cfg.FrancisHostPSK = []byte("bootstrap-psk-that-is-long-enough")
|
||||
|
||||
_, err := NewActorStateStore(NewActorsOpts{
|
||||
EnvConfig: cfg,
|
||||
InstanceID: "ee05c3eb-8129-47a6-a1c7-849998b6f876",
|
||||
DB: newDB(t),
|
||||
})
|
||||
require.ErrorIs(t, err, ErrRemoteFrancisRuntime)
|
||||
})
|
||||
}
|
||||
|
||||
// TestNewActorsBackupProvider covers the provider the export and import use to back up and restore the actor host's data.
|
||||
// It builds the provider from the same options the actor host uses, so a mismatch between those options and the concrete provider would otherwise only surface at runtime, when an export or import is attempted.
|
||||
func TestNewActorsBackupProvider(t *testing.T) {
|
||||
@@ -70,3 +228,26 @@ func TestNewActorsBackupProvider(t *testing.T) {
|
||||
err = provider.Restore(t.Context(), bytes.NewReader(buf.Bytes()))
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
// testCAPEM returns a self-signed CA certificate in PEM form, standing in for the cluster CA an operator would pin with FRANCIS_CA
|
||||
func testCAPEM(t *testing.T) []byte {
|
||||
t.Helper()
|
||||
|
||||
pub, priv, err := ed25519.GenerateKey(rand.Reader)
|
||||
require.NoError(t, err)
|
||||
|
||||
tmpl := &x509.Certificate{
|
||||
SerialNumber: big.NewInt(1),
|
||||
Subject: pkix.Name{CommonName: "test-cluster-ca"},
|
||||
NotBefore: time.Now().Add(-time.Hour),
|
||||
NotAfter: time.Now().Add(time.Hour),
|
||||
IsCA: true,
|
||||
KeyUsage: x509.KeyUsageCertSign,
|
||||
BasicConstraintsValid: true,
|
||||
}
|
||||
|
||||
der, err := x509.CreateCertificate(rand.Reader, tmpl, tmpl, pub, priv)
|
||||
require.NoError(t, err)
|
||||
|
||||
return pem.EncodeToMemory(&pem.Block{Type: "CERTIFICATE", Bytes: der})
|
||||
}
|
||||
|
||||
@@ -0,0 +1,78 @@
|
||||
package bootstrap
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"time"
|
||||
|
||||
francishost "github.com/italypaleale/francis/host"
|
||||
"github.com/italypaleale/francis/host/remote"
|
||||
|
||||
"github.com/pocket-id/pocket-id/backend/internal/common"
|
||||
)
|
||||
|
||||
// actorClientConnectTimeout bounds how long a CLI command waits to join the cluster
|
||||
// Francis reconnects to the runtime indefinitely, which is right for the server but would leave a command hanging against an unreachable runtime
|
||||
const actorClientConnectTimeout = 30 * time.Second
|
||||
|
||||
// WithActorClient connects to the standalone Francis runtime, calls fn once the connection is live, and disconnects before returning.
|
||||
// It's meant for CLI commands, which have no actor host of their own: the client joins the cluster only for the duration of fn, and hosts no actor while connected, so the runtime never places an actor on it.
|
||||
// It requires FRANCIS_HOST to point to a standalone runtime, and returns ErrEmbeddedFrancisRuntime otherwise, since an embedded runtime is reached through the database instead.
|
||||
func WithActorClient(parentCtx context.Context, envConfig *common.EnvConfigSchema, fn func(ctx context.Context, client francishost.Host) error) error {
|
||||
if envConfig.HasEmbeddedFrancisRuntime() {
|
||||
return ErrEmbeddedFrancisRuntime
|
||||
}
|
||||
|
||||
log := slog.Default().With("scope", "actor-client")
|
||||
|
||||
// The client hosts no actor, so it advertises no address of its own and binds nothing
|
||||
// The short grace period keeps a command from lingering on the way out, since there are no actors to drain
|
||||
client, err := remote.NewHost(append(
|
||||
remoteConnectionOptions(envConfig, log),
|
||||
remote.WithClientOnly(),
|
||||
remote.WithShutdownGracePeriod(2*time.Second),
|
||||
)...)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to create the actor client: %w", err)
|
||||
}
|
||||
|
||||
ctx, cancel := context.WithCancel(parentCtx)
|
||||
defer cancel()
|
||||
|
||||
runErrCh := make(chan error, 1)
|
||||
go func() {
|
||||
runErrCh <- client.Run(ctx)
|
||||
}()
|
||||
|
||||
// Every operation travels on the runtime session, so nothing can run before the client has joined the cluster
|
||||
connectCtx, connectCancel := context.WithTimeout(ctx, actorClientConnectTimeout)
|
||||
defer connectCancel()
|
||||
|
||||
select {
|
||||
case <-client.Ready():
|
||||
case runErr := <-runErrCh:
|
||||
return fmt.Errorf("failed to connect to the Francis runtime: %w", runErr)
|
||||
case <-connectCtx.Done():
|
||||
cancel()
|
||||
<-runErrCh
|
||||
return fmt.Errorf("timed out connecting to the Francis runtime after %v", actorClientConnectTimeout)
|
||||
}
|
||||
|
||||
fnErr := fn(ctx, client)
|
||||
|
||||
// Leave the cluster before returning, so the runtime drops the registration instead of waiting for the health check to lapse
|
||||
cancel()
|
||||
runErr := <-runErrCh
|
||||
|
||||
// The error from fn is the one the caller asked for, and a canceled run is just the disconnect we asked for
|
||||
switch {
|
||||
case fnErr != nil:
|
||||
return fnErr
|
||||
case runErr != nil && !errors.Is(runErr, context.Canceled):
|
||||
return fmt.Errorf("error disconnecting from the Francis runtime: %w", runErr)
|
||||
default:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
@@ -10,7 +10,7 @@ import (
|
||||
_ "github.com/golang-migrate/migrate/v4/source/file"
|
||||
|
||||
"github.com/italypaleale/francis/components"
|
||||
"github.com/italypaleale/francis/host/local"
|
||||
francishost "github.com/italypaleale/francis/host"
|
||||
"github.com/italypaleale/go-kit/servicerunner"
|
||||
"gorm.io/gorm"
|
||||
|
||||
@@ -137,7 +137,7 @@ func Bootstrap(ctx context.Context) error {
|
||||
}
|
||||
|
||||
// actorsRunServiceFn wraps the actor host's Run method in a background service and returns a "ready" signal that other services can wait on
|
||||
func actorsRunServiceFn(actors *local.Host) (servicerunner.Service, *servicerunner.Ready) {
|
||||
func actorsRunServiceFn(actors francishost.Host) (servicerunner.Service, *servicerunner.Ready) {
|
||||
actorsReady := servicerunner.NewReady()
|
||||
fn := func(ctx context.Context) error {
|
||||
runErrCh := make(chan error, 1)
|
||||
|
||||
@@ -5,7 +5,7 @@ import (
|
||||
"fmt"
|
||||
"net/http"
|
||||
|
||||
"github.com/italypaleale/francis/host/local"
|
||||
francishost "github.com/italypaleale/francis/host"
|
||||
"github.com/pocket-id/pocket-id/backend/internal/api"
|
||||
"github.com/pocket-id/pocket-id/backend/internal/apikey"
|
||||
"github.com/pocket-id/pocket-id/backend/internal/appconfig"
|
||||
@@ -52,7 +52,7 @@ type services struct {
|
||||
emailVerificationModule *emailverification.Module
|
||||
apiModule *api.Module
|
||||
environmentModule *environment.Module
|
||||
actors *local.Host
|
||||
actors francishost.Host
|
||||
}
|
||||
|
||||
// Initializes all services
|
||||
@@ -60,7 +60,7 @@ func initServices(
|
||||
ctx context.Context,
|
||||
db *gorm.DB,
|
||||
instanceID string,
|
||||
actors *local.Host,
|
||||
actors francishost.Host,
|
||||
httpClient *http.Client,
|
||||
imageExtensions map[string]string,
|
||||
fileStorage storage.FileStorage,
|
||||
|
||||
Reference in New Issue
Block a user