mirror of
https://github.com/pocket-id/pocket-id.git
synced 2026-09-29 14:29:04 +02:00
feat: add FRANCIS_HOST to connect to a standalone Francis runtime
FRANCIS_HOST decides where the Francis actor runtime lives. When set to "embedded" (the default), Pocket ID starts the runtime inside its own process. 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. Because when using a remote runtime, it's likewise not possible to enforce a single instance of Pocket ID is running at once, the env vars currently have the `EXPERIMENTAL_` prefix, are **undocumented**, and show a warning if used. Notes: - Connecting to a standalone runtime also needs FRANCIS_HOST_PSK or FRANCIS_HOST_JWT_FILE, and optionally (but recommended) FRANCIS_CA. - When connecting to a remote runtime, exporting Pocket ID data does not include the actor state, which will need to be backed up and restored separately
This commit is contained in:
@@ -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,14 +44,50 @@ type NewActorsOpts struct {
|
||||
FileStorage storage.FileStorage
|
||||
}
|
||||
|
||||
func NewActors(o NewActorsOpts) (*local.Host, map[string]*ratelimit.RateLimitService, error) {
|
||||
log := slog.Default()
|
||||
func NewActors(o NewActorsOpts) (h francishost.Host, rateLimitServices map[string]*ratelimit.RateLimitService, err error) {
|
||||
log := slog.Default().With("scope", "actor-host")
|
||||
|
||||
// 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
|
||||
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)
|
||||
}
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Add all cron jobs
|
||||
err = o.registerCronJobs(h)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Add the rate limiters
|
||||
rateLimiters, err := o.registerRateLimiters(h)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Bind a service for each rate limiter so the middleware can invoke them
|
||||
rateLimitServices = make(map[string]*ratelimit.RateLimitService, len(rateLimiters))
|
||||
for name, rl := range rateLimiters {
|
||||
rateLimitServices[name] = rl.Service(h.Service())
|
||||
}
|
||||
|
||||
return h, rateLimitServices, nil
|
||||
}
|
||||
|
||||
// newEmbeddedHost creates the actor host that runs the Francis runtime inside the Pocket ID process
|
||||
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, nil, fmt.Errorf("failed to derive PSK: %w", err)
|
||||
return nil, fmt.Errorf("failed to derive PSK: %w", err)
|
||||
}
|
||||
|
||||
// Derive the cluster host limit from the HA setting
|
||||
@@ -55,7 +101,7 @@ func NewActors(o NewActorsOpts) (*local.Host, map[string]*ratelimit.RateLimitSer
|
||||
// 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.WithLogger(log),
|
||||
local.WithRuntimePSKs(psk),
|
||||
local.WithShutdownGracePeriod(10 * time.Second),
|
||||
local.WithMaxHosts(maxHosts),
|
||||
@@ -76,35 +122,59 @@ func NewActors(o NewActorsOpts) (*local.Host, map[string]*ratelimit.RateLimitSer
|
||||
// Add the database connection
|
||||
providerOpt, err := o.getProviderOption()
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
return 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)
|
||||
return nil, fmt.Errorf("failed to create actor host: %w", err)
|
||||
}
|
||||
|
||||
// Add all cron jobs
|
||||
err = o.registerCronJobs(h)
|
||||
return h, nil
|
||||
}
|
||||
|
||||
// newRemoteHost creates the actor host that connects to a standalone Francis runtime
|
||||
func (o *NewActorsOpts) newRemoteHost(log *slog.Logger) (*remote.Host, error) {
|
||||
opts := append(
|
||||
remoteConnectionOptions(o.EnvConfig, log),
|
||||
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, nil, err
|
||||
return nil, fmt.Errorf("failed to create remote actor host: %w", err)
|
||||
}
|
||||
|
||||
// Add the rate limiters
|
||||
rateLimiters, err := o.registerRateLimiters(h)
|
||||
if err != nil {
|
||||
return nil, nil, 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()...),
|
||||
}
|
||||
|
||||
// Bind a service for each rate limiter so the middleware can invoke them
|
||||
rateLimitServices := make(map[string]*ratelimit.RateLimitService, len(rateLimiters))
|
||||
for name, rl := range rateLimiters {
|
||||
rateLimitServices[name] = rl.Service(h.Service())
|
||||
// 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))
|
||||
}
|
||||
|
||||
return h, rateLimitServices, nil
|
||||
// Pinning the cluster CA lets Pocket ID verify the runtime on its very first connection
|
||||
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
|
||||
@@ -117,7 +187,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 +309,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 +346,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,150 @@ 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.SetFrancisAddresses([]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 file": func(cfg *common.EnvConfigSchema) { cfg.FrancisHostJWTFile = jwtFile },
|
||||
} {
|
||||
t.Run(name, func(t *testing.T) {
|
||||
cfg := baseConfig(t)
|
||||
cfg.SetFrancisAddresses([]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.SetFrancisAddresses([]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.SetFrancisAddresses([]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.SetFrancisAddresses([]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.SetFrancisAddresses([]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 +227,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 is the timeout for 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 err = <-runErrCh:
|
||||
return fmt.Errorf("failed to connect to the Francis runtime: %w", err)
|
||||
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()
|
||||
err = <-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 err != nil && !errors.Is(err, context.Canceled):
|
||||
return fmt.Errorf("error disconnecting from the Francis runtime: %w", err)
|
||||
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