/
githubmirror
/
photoprism
Обзор
Документация
Войти
/
githubmirror
/
photoprism
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
develop
internal/service/cluster/node/bootstrap.go
725 строк
22 KB
Michael Mayer
Cluster: Re-register node config on restart via OAuth credentials
13 июн 2026, 15:54
13 июн 2026, 15:54
974321b
Код
Авторство
О чём код?
package node import ( "context" "encoding/base64" "encoding/json" "errors" "fmt" "io" "net" "net/http" "net/url" "os" "path/filepath" "strconv" "strings" "time" clusterjwt "github.com/photoprism/photoprism/internal/auth/jwt" "github.com/photoprism/photoprism/internal/config" "github.com/photoprism/photoprism/internal/event" "github.com/photoprism/photoprism/internal/service/cluster" "github.com/photoprism/photoprism/pkg/clean" "github.com/photoprism/photoprism/pkg/fs" "github.com/photoprism/photoprism/pkg/http/dns" "github.com/photoprism/photoprism/pkg/http/scheme" "github.com/photoprism/photoprism/pkg/rnd" ) var log = event.Log const bootstrapOAuthScope = "cluster" // init registers the cluster node bootstrap extension so it runs before the // database connection is established. func init() { // Register early so this can adjust DB settings before connectDb(). config.Register(config.StageBoot, "cluster-node", InitConfig, nil) } // InitConfig performs node bootstrap: optional registration with the Portal // and theme installation. Runs early during config.Init(). func InitConfig(c *config.Config) error { role := c.NodeRole() // Skip on portal nodes and unknown node types. if c.Portal() || (role != cluster.RoleInstance && role != cluster.RoleService) { log.Debugf("config: skipping cluster bootstrap for %s", clean.Log(role)) return nil } // Auto-join the cluster and sync the theme when enabled and configured. if cluster.BootstrapAutoJoinEnabled || cluster.BootstrapAutoThemeEnabled { bootstrapClusterNode(c) } // Derive the OIDC RP credentials from the node client (PHOTOPRISM_CLUSTER_OIDC), // independent of the auto-join/theme toggles, so a registered instance re-wires // the OIDC RP on every restart. resolveNodeOIDCClient(c) // Log cluster UUID. if uuid := c.ClusterUUID(); uuid != "" { log.Infof("cluster: UUID %s", clean.Log(uuid)) } return nil } // bootstrapClusterNode registers the node with the configured Portal and installs // its theme, honoring the auto-join/theme toggles. All failures are non-fatal so // the node still boots; it is a no-op when no Portal URL or join token is set. func bootstrapClusterNode(c *config.Config) { portalURL := strings.TrimSpace(c.PortalUrl()) joinToken := strings.TrimSpace(c.JoinToken()) // Proceed with either a join token (first join) or the node's own OAuth // credentials (already joined). The credential-only case lets a plain restart // re-send the idempotent registration so a changed declarative config reaches // the Portal after the join token has been removed from the environment. hasNodeCreds := strings.TrimSpace(c.NodeClientID()) != "" && strings.TrimSpace(c.NodeClientSecret()) != "" if portalURL == "" || (joinToken == "" && !hasNodeCreds) { log.Debugf("cluster: no bootstrap configuration found") return } log.Debugf("config: attempting to join the configured cluster") u, err := url.Parse(portalURL) if err != nil || u.Scheme == "" || u.Host == "" { log.Warnf("cluster: invalid portal URL %s", clean.Log(portalURL)) return } // Register with retry policy. var registerResp *cluster.RegisterResponse if cluster.BootstrapAutoJoinEnabled { if registerResp, err = registerWithPortal(c, u, joinToken); err != nil { log.Warnf("config: failed to join the configured cluster (%s)", clean.Error(err)) } } // Pull theme if missing or outdated, and activate it when present. if cluster.BootstrapAutoThemeEnabled { if err = syncNodeTheme(c, u, registerResp); err != nil { // Theme install failures are non-critical; log at debug to avoid noise. log.Debugf("cluster: theme download skipped (%s)", clean.Error(err)) } activateNodeThemeIfPresent(c) } } // resolveNodeOIDCClient derives the instance's OIDC RP credentials from the node // client credentials when PHOTOPRISM_CLUSTER_OIDC is enabled, so a Portal-fronted // instance logs in on first boot without injecting PHOTOPRISM_OIDC_*. An explicit // OIDC client id wins; the issuer defaults to the instance's own origin when unset; // the secret is read file-first so a rotation propagates on the next start. func resolveNodeOIDCClient(c *config.Config) { if c == nil || !c.ClusterOIDC() { return } // Explicit OIDC client credentials (a different IdP, or a hand-issued client) // win unchanged. if strings.TrimSpace(c.OIDCClient()) != "" { return } if c.NodeRole() != cluster.RoleInstance { log.Warnf("cluster: ignoring cluster OIDC because this node is not an instance") return } id := strings.TrimSpace(c.NodeClientID()) secret := strings.TrimSpace(c.NodeClientSecret()) if id == "" || secret == "" { log.Warnf("cluster: cannot derive the OIDC client from node credentials yet (node not registered)") return } // Default the OIDC issuer to the instance's own origin root (the shared-domain // Portal OP) when unset, so enabling cluster OIDC is the only configuration an // instance needs. An explicit PHOTOPRISM_OIDC_URI (e.g. a subdomain-isolated // Portal) is respected. if c.OIDCUri().Host == "" { if issuer := scheme.OriginURL(c.SiteUrl()); issuer != "" { c.SetOIDCUri(issuer) } } c.SetOIDCClient(id) c.SetOIDCSecret(secret) log.Infof("cluster: OIDC login configured via the Portal using node client %s", clean.Log(id)) } // newHTTPClient returns a short-lived HTTP client configured with the provided // timeout. It is intentionally lightweight to avoid leaking transports between // bootstrap attempts. func newHTTPClient(timeout time.Duration) *http.Client { // TODO: Consider reusing a shared *http.Transport with sane defaults and enabling // proxy support explicitly if required. For now, rely on net/http defaults and // the HTTPS_PROXY set in config.Init(). return &http.Client{Timeout: timeout} } // registerWithPortal attempts to register the node with the Portal, retrying on // transient errors up to the configured limits. Successful registrations update // local configuration and prime JWKS credentials. func registerWithPortal(c *config.Config, portal *url.URL, token string) (*cluster.RegisterResponse, error) { maxAttempts := cluster.BootstrapRegisterMaxAttempts delay := cluster.BootstrapRegisterRetryDelay timeout := cluster.BootstrapRegisterTimeout endpoint := *portal endpoint.Path = strings.TrimRight(endpoint.Path, "/") + "/api/v1/cluster/nodes/register" // Let the configuration decide if credentials are missing (MySQL with no effective name/user/password). wantRotateDatabase := c.ShouldAutoRotateDatabase() payload := buildRegisterPayload(c) if wantRotateDatabase { // Align with API: request database rotation/creation on (re)register. payload.RotateDatabase = true } authToken, err := registerAuthToken(c, portal, token) if err != nil { return nil, err } bodyBytes, _ := json.Marshal(payload) for attempt := 1; attempt <= maxAttempts; attempt++ { req, _ := http.NewRequest(http.MethodPost, endpoint.String(), strings.NewReader(string(bodyBytes))) req.Header.Set("Authorization", "Bearer "+authToken) req.Header.Set("Content-Type", "application/json") req.Header.Set("Accept", "application/json") // Endpoint is derived from the configured Portal URL. resp, err := newHTTPClient(timeout).Do(req) //nolint:gosec if err != nil { if attempt < maxAttempts { log.Debugf("cluster: join attempt %d/%d failed with %s", attempt, maxAttempts, clean.Error(err)) time.Sleep(delay) continue } return nil, err } retry, registerResp, err := func() (bool, *cluster.RegisterResponse, error) { defer func() { if closeErr := resp.Body.Close(); closeErr != nil { log.Debugf("cluster: %s (close registration response body)", clean.Error(closeErr)) } }() switch resp.StatusCode { case http.StatusOK, http.StatusCreated: var r cluster.RegisterResponse dec := json.NewDecoder(resp.Body) if err = dec.Decode(&r); err != nil { return false, nil, err } if err = persistRegistration(c, &r, wantRotateDatabase); err != nil { return false, nil, err } primeJWKS(c, r.JWKSUrl) if resp.StatusCode == http.StatusCreated { log.Infof("config: successfully joined cluster as instance %s (%d)", clean.LogQuote(r.Node.Name), resp.StatusCode) } else { log.Infof("cluster: membership confirmed") } return false, &r, nil case http.StatusUnauthorized, http.StatusForbidden, http.StatusNotFound: // Terminal errors (no retry). 404 likely indicates a Portal without cluster endpoints. return false, nil, errors.New(resp.Status) case http.StatusTooManyRequests: if attempt < maxAttempts { log.Debugf("cluster: join attempt %d/%d rate limited by portal", attempt, maxAttempts) return true, nil, nil } return false, nil, errors.New(resp.Status) case http.StatusConflict, http.StatusBadRequest: // Do not retry on 400/409 per spec intent. return false, nil, errors.New(resp.Status) default: if attempt < maxAttempts { log.Debugf("cluster: join attempt %d/%d failed with status %s", attempt, maxAttempts, resp.Status) // TODO: Consider exponential backoff with jitter instead of constant delay. return true, nil, nil } return false, nil, errors.New(resp.Status) } }() switch { case err != nil: return nil, err case registerResp != nil: return registerResp, nil case retry: time.Sleep(delay) continue } } return nil, nil } // registerAuthToken returns the bearer token used for register requests. // Existing-node mutations use an OAuth access token, while first-time joins // use the configured join token when no node credentials exist yet. func registerAuthToken(c *config.Config, portal *url.URL, joinToken string) (string, error) { if c == nil || portal == nil { return "", fmt.Errorf("invalid cluster bootstrap config") } if id, secret := strings.TrimSpace(c.NodeClientID()), strings.TrimSpace(c.NodeClientSecret()); id != "" && secret != "" { token, err := oauthAccessToken(portal, id, secret, bootstrapOAuthScope) if err != nil { return "", fmt.Errorf("portal access token request failed: %w", err) } return token, nil } if token := strings.TrimSpace(joinToken); token != "" { return token, nil } return "", fmt.Errorf("missing join token and node client credentials") } // defaultClusterDomain returns the configured cluster domain or, if absent, // attempts to derive it from the Portal URL by stripping common prefixes. func defaultClusterDomain(c *config.Config) string { if c == nil { return "" } domain := strings.TrimSpace(c.ClusterDomain()) if domain != "" { return strings.Trim(domain, ".") } portalURL := strings.TrimSpace(c.PortalUrl()) if portalURL == "" { return "" } u, err := url.Parse(portalURL) if err != nil { return "" } host := strings.Trim(u.Hostname(), ".") if host == "" { return "" } // Strip common prefixes like portal.<domain>. if dns.IsLoopbackHost(host) { return "" } if ip := net.ParseIP(host); ip != nil { // Prefer DNS domains over raw IP addresses; leave empty so caller can decide. return "" } if strings.HasPrefix(host, "portal.") && len(host) > len("portal.") { return strings.TrimPrefix(host, "portal.") } return host } // defaultNodeURL builds https://<name>.<domain> using sanitized labels. func defaultNodeURL(name, domain string) string { name = clean.TypeLowerDash(strings.TrimSpace(name)) domain = strings.Trim(strings.ToLower(domain), ".") if name == "" || domain == "" { return "" } return fmt.Sprintf("https://%s.%s", name, domain) } // buildRegisterPayload builds a registration payload with stable defaults so // all registration code paths report consistent node metadata. func buildRegisterPayload(c *config.Config) cluster.RegisterRequest { payload := cluster.RegisterRequest{ NodeName: c.NodeName(), NodeUUID: c.NodeUUID(), NodeRole: c.NodeRole(), AdvertiseUrl: c.AdvertiseUrl(), AppName: clean.TypeUnicode(c.About()), AppVersion: clean.TypeUnicode(c.Version()), Theme: clean.TypeUnicode(c.NodeThemeVersion()), } // Report a human-friendly DisplayName from the operator's configured branding via // Config.SiteName (SITE_NAME, then the raw AppName, then SiteTitle). It is empty // for an unbranded instance, so the Portal falls back to the node Name slug, and // it ignores the product Name default so an unbranded Pro node does not look // configured. SiteCaption is intentionally excluded: Plus/Pro default it to the // shared marketing description, so it is not a per-instance label. payload.DisplayName = c.SiteName() // Auto-derive Advertise/Site URLs from node name and cluster domain when not configured. if domain := strings.TrimSpace(defaultClusterDomain(c)); domain != "" { if payload.NodeName == "" { payload.NodeName = c.NodeName() } if payload.AdvertiseUrl == "" { if u := defaultNodeURL(payload.NodeName, domain); u != "" { payload.AdvertiseUrl = u } } if payload.SiteUrl == "" && payload.AdvertiseUrl != "" { payload.SiteUrl = payload.AdvertiseUrl } } // Include SiteUrl whenever configured; the server normalizes duplicates if needed. if su := c.SiteUrl(); su != "" { payload.SiteUrl = su } // Declare the instance's group-based admission config so it takes effect on // first boot and every re-registration; an admin override on the Portal // still wins (see the registry's group-config source precedence). payload.AllowGroups = c.ClusterAllowGroups() payload.AllowGroupRoles = c.ClusterAllowGroupRoles() payload.GroupsFullView = c.ClusterGroupsFullView() return payload } // persistRegistration stores registration responses through the config package // so cluster option writes and secret-file persistence stay in one place. func persistRegistration(c *config.Config, r *cluster.RegisterResponse, wantRotateDatabase bool) error { updates := cluster.OptionsUpdate{} // Persist ClusterUUID from portal response if provided. if rnd.IsUUID(r.UUID) { updates.SetClusterUUID(r.UUID) } if cidr := strings.TrimSpace(r.ClusterCIDR); cidr != "" { updates.SetClusterCIDR(cidr) } // Always persist NodeClientID (client UID) from response for future OAuth token requests. if r.Node.ClientID != "" { updates.SetNodeClientID(r.Node.ClientID) } // Persist node client secret only if missing locally and provided by server. if r.Secrets != nil && r.Secrets.ClientSecret != "" && c.NodeClientSecret() == "" { if _, err := c.SaveNodeClientSecret(r.Secrets.ClientSecret); err != nil { return fmt.Errorf("failed to persist node client secret: %w", err) } } if jwksUrl := strings.TrimSpace(r.JWKSUrl); jwksUrl != "" { updates.SetJWKSUrl(jwksUrl) c.SetJWKSUrl(jwksUrl) } // Apply the validating setter first and persist only what it accepted, so // a malformed URL from a misbehaving Portal is never written to options.yml. if loginUrl := strings.TrimSpace(r.PortalLoginUrl); loginUrl != "" { c.SetPortalLoginUrl(loginUrl) if v := c.PortalLoginUrl(); v != "" { updates.SetPortalLoginUrl(v) } } // Persist NodeUUID from portal response if provided and not set locally. if r.Node.UUID != "" && c.NodeUUID() == "" { updates.SetNodeUUID(r.Node.UUID) } // Persist DB settings only if rotation was requested and driver is MySQL/MariaDB // and local DB not configured (as checked before calling). if wantRotateDatabase { if r.Database.DSN != "" { updates.SetDatabaseDriver(r.Database.Driver) updates.SetDatabaseDSN(r.Database.DSN) } else if r.Database.Name != "" && r.Database.User != "" && r.Database.Password != "" { server := r.Database.Host if r.Database.Port > 0 { server = net.JoinHostPort(r.Database.Host, strconv.Itoa(r.Database.Port)) } updates.SetDatabaseDriver(r.Database.Driver) updates.SetDatabaseServer(server) updates.SetDatabaseName(r.Database.Name) updates.SetDatabaseUser(r.Database.User) updates.SetDatabasePassword(r.Database.Password) } } if updates.IsZero() { return nil } wrote, err := c.SaveClusterOptionsUpdate(updates) if err != nil { return err } if wrote && updates.HasDatabaseUpdate() { log.Infof("config: applied portal database settings; restart required to connect with new credentials") } return nil } // primeJWKS eagerly fetches the Portal JWKS so that subsequent token // verification does not incur network latency during critical operations. func primeJWKS(c *config.Config, url string) { if c == nil { return } url = strings.TrimSpace(url) if url == "" { return } verifier := clusterjwt.NewVerifier(c) ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) defer cancel() if err := verifier.Prime(ctx, url); err != nil { log.Debugf("auth: jwks prime skipped (%s)", clean.Error(err)) } } // syncNodeTheme downloads or refreshes the Portal-provided theme in the instance-specific // theme directory (NodeThemePath) when the local version is missing or differs from the portal version. func syncNodeTheme(c *config.Config, portal *url.URL, registerResp *cluster.RegisterResponse) error { themeDir := c.NodeThemePath() localVersion := strings.TrimSpace(c.NodeThemeVersion()) hasAppJS := fs.FileExists(filepath.Join(themeDir, fs.AppJsFile)) portalVersion := "" if registerResp != nil { portalVersion = clean.TypeUnicode(registerResp.Theme) } shouldProbe := registerResp == nil needsDownload := false requiresOverwrite := false switch { case portalVersion != "": switch { case !hasAppJS: log.Infof("theme: %s not installed yet; scheduling download", clean.Log(portalVersion)) needsDownload = true case localVersion != portalVersion: log.Infof("theme: update detected (local %s, portal %s); scheduling download", clean.Log(localVersion), clean.Log(portalVersion)) needsDownload = true requiresOverwrite = true default: log.Infof("theme: version %s already installed", clean.Log(localVersion)) } case shouldProbe: // Registration failed or was skipped; attempt to obtain the theme when missing. needsDownload = !hasAppJS || localVersion == "" if needsDownload { log.Infof("theme: probing portal because local bundle is missing") } default: // Portal responded but has no theme configured; keep existing node theme. log.Infof("cluster: portal did not advertise a theme; skipping download") return nil } if !needsDownload { return nil } // Acquire OAuth bearer via client credentials; skip when credentials are unavailable. bearer := "" var tokenErr error if id, secret := strings.TrimSpace(c.NodeClientID()), strings.TrimSpace(c.NodeClientSecret()); id != "" && secret != "" { if t, err := oauthAccessToken(portal, id, secret, bootstrapOAuthScope); err != nil { tokenErr = err log.Infof("config: portal access token request failed (%s)", clean.Error(err)) } else { bearer = t } } if bearer == "" { if tokenErr != nil { log.Infof("theme: sync skipped because portal access token request failed (%s)", clean.Error(tokenErr)) } } if bearer == "" { log.Infof("theme: sync skipped because no portal credentials are available yet") return nil } endpoint := *portal endpoint.Path = strings.TrimRight(endpoint.Path, "/") + "/api/v1/cluster/theme" req, _ := http.NewRequest(http.MethodGet, endpoint.String(), nil) req.Header.Set("Authorization", "Bearer "+bearer) req.Header.Set("Accept", "application/zip") // Endpoint is derived from the configured Portal URL. resp, err := newHTTPClient(cluster.BootstrapRegisterTimeout).Do(req) //nolint:gosec if err != nil { return err } defer func() { if closeErr := resp.Body.Close(); closeErr != nil { log.Debugf("theme: %s (close theme response body)", clean.Error(closeErr)) } }() switch resp.StatusCode { case http.StatusOK: if err = fs.MkdirAll(c.TempPath()); err != nil { return err } zipName := filepath.Join(c.TempPath(), "cluster-theme.zip") var out *os.File if out, err = os.Create(zipName); err != nil { //nolint:gosec return err } if _, err = io.Copy(out, resp.Body); err != nil { _ = out.Close() return err } _ = out.Close() if requiresOverwrite && fs.PathExists(themeDir) { if err = os.RemoveAll(themeDir); err != nil { return err } } if err = fs.MkdirAll(themeDir); err != nil { return err } _, _, unzipErr := fs.Unzip(zipName, themeDir, 32*fs.MB, 512*fs.MB) return unzipErr case http.StatusNotFound: return nil case http.StatusUnauthorized, http.StatusForbidden: return errors.New(resp.Status) default: return errors.New(resp.Status) } } // activateNodeThemeIfPresent switches the active theme path to the instance-specific // NodeThemePath directory when a valid cluster-managed theme bundle is available. func activateNodeThemeIfPresent(c *config.Config) { if c == nil { return } // If NodeThemePath() does not exist or does not contain an app.js file, // NodeThemeVersion() returns an empty string. No additional checks required. if c.NodeThemeVersion() == "" { return } // nodeDir is already clean, because filepath.Join() returns it that way. nodeDir := c.NodeThemePath() // Return is theme is already activated. if filepath.Clean(c.ThemePath()) == nodeDir { return } // Activate cluster theme. c.SetThemePath(nodeDir) // Report activation. log.Debugf("config: activated portal theme from %s", clean.Log(nodeDir)) } // oauthAccessToken requests an OAuth access token via client_credentials using Basic auth. func oauthAccessToken(portal *url.URL, clientID, clientSecret, scope string) (string, error) { if portal == nil { return "", fmt.Errorf("invalid portal url") } tokenURL := *portal tokenURL.Path = strings.TrimRight(tokenURL.Path, "/") + "/api/v1/oauth/token" form := url.Values{} form.Set("grant_type", "client_credentials") form.Set("scope", clean.Scope(scope)) req, _ := http.NewRequest(http.MethodPost, tokenURL.String(), strings.NewReader(form.Encode())) req.Header.Set("Content-Type", "application/x-www-form-urlencoded") req.Header.Set("Accept", "application/json") // Basic auth for client credentials basic := base64.StdEncoding.EncodeToString([]byte(clientID + ":" + clientSecret)) req.Header.Set("Authorization", "Basic "+basic) // Endpoint is derived from the configured Portal URL. resp, err := newHTTPClient(cluster.BootstrapRegisterTimeout).Do(req) //nolint:gosec if err != nil { return "", err } defer func() { if closeErr := resp.Body.Close(); closeErr != nil { log.Debugf("cluster: %s (close token response body)", clean.Error(closeErr)) } }() if resp.StatusCode != http.StatusOK { return "", fmt.Errorf("%s", resp.Status) } var tok map[string]any dec := json.NewDecoder(resp.Body) if err = dec.Decode(&tok); err != nil { return "", err } accessToken, _ := tok["access_token"].(string) if accessToken == "" { return "", fmt.Errorf("empty access_token") } return accessToken, nil }