This is an automated email from the ASF dual-hosted git repository.

nickva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/couchdb.git


The following commit(s) were added to refs/heads/main by this push:
     new c4fd3366c Supervisor-level replicator auth plugin callbacks
c4fd3366c is described below

commit c4fd3366c3fcd073a92e8d14027fdd3fa7ad7c10
Author: Nick Vatamaniuc <[email protected]>
AuthorDate: Mon Jul 20 17:54:57 2026 -0400

    Supervisor-level replicator auth plugin callbacks
    
    As discussed in the comments of the new IAM plugin in [1], it would be nice 
to
    have a per-plugin supervisor-level context for each plugin. So, for example,
    they can create some kind of a ETS cache table for their tokens. This is 
what
    we implement here. The API is simple:
    
      * `sup_initialize() -> Ctx`
      * `sup_cleanup(Ctx) -> ok`
    
    The APIs are optional to implement. The name pattern mirror the regular 
plugin
    API names: `initialize(...) -> {ok, ..., Ctx}` and `cleanup(Ctx) -> ok`
    
    [1] https://github.com/apache/couchdb/pull/6069
---
 src/couch_replicator/src/couch_replicator_auth.erl | 155 ++++++++++++++++++++-
 src/couch_replicator/src/couch_replicator_sup.erl  |   3 +-
 2 files changed, 156 insertions(+), 2 deletions(-)

diff --git a/src/couch_replicator/src/couch_replicator_auth.erl 
b/src/couch_replicator/src/couch_replicator_auth.erl
index 712a771a4..98d7d8404 100644
--- a/src/couch_replicator/src/couch_replicator_auth.erl
+++ b/src/couch_replicator/src/couch_replicator_auth.erl
@@ -12,6 +12,9 @@
 
 -module(couch_replicator_auth).
 
+-behaviour(gen_server).
+-behaviour(config_listener).
+
 -export([
     initialize/1,
     update_headers/2,
@@ -19,12 +22,35 @@
     cleanup/1
 ]).
 
+-export([
+    start_link/0,
+    reconcile/0,
+    context/0
+]).
+
+-export([
+    init/1,
+    handle_call/3,
+    handle_cast/2,
+    handle_info/2,
+    terminate/2
+]).
+
+-export([
+    handle_config_change/5,
+    handle_config_terminate/3
+]).
+
 -include_lib("couch_replicator/include/couch_replicator_api_wrap.hrl").
 
 -type headers() :: [{string(), string()}].
 -type code() :: non_neg_integer().
 
 -define(DEFAULT_PLUGINS, 
"couch_replicator_auth_session,couch_replicator_auth_noop").
+-define(RELISTEN_DELAY, 5000).
+
+% Plugin module map to whatever that plugin's sup_initialize/0 returned.
+-record(st, {ctx = #{}}).
 
 % Behavior API
 
@@ -44,6 +70,21 @@
 
 -callback cleanup(term()) -> ok.
 
+% Optional. Node-wide setup and teardown, A plugin that needs state shared by
+% every replication job (say, a token ets cache) starts it in sup_initialize/0
+% and returns an opaque context. That context is handed back to sup_cleanup/1
+% when the plugin is removed from the `[replicator] auth_plugins` config or the
+% replicator shuts down.
+%
+% Anything linked here is linked to the calling process, which is
+% couch_replicator_auth gen_server.
+%
+-callback sup_initialize() -> {ok, term()} | ignore | {error, term()}.
+
+-callback sup_cleanup(term()) -> ok.
+
+-optional_callbacks([sup_initialize/0, sup_cleanup/1]).
+
 % Main API
 
 -spec initialize(#httpdb{}) -> {ok, #httpdb{}} | {error, term()}.
@@ -72,7 +113,119 @@ cleanup(#httpdb{auth_context = {Module, Context}} = 
HttpDb) ->
     ok = Module:cleanup(Context),
     HttpDb#httpdb{auth_context = nil}.
 
-% Private helper functions
+% Plugin manager to manage the optional sub_initialize/0 contexts.
+%
+
+start_link() ->
+    gen_server:start_link({local, ?MODULE}, ?MODULE, [], []).
+
+-spec reconcile() -> ok.
+reconcile() ->
+    gen_server:call(?MODULE, reconcile, infinity).
+
+% The plugin modules currently set up.
+-spec context() -> [atom()].
+context() ->
+    gen_server:call(?MODULE, context, infinity).
+
+init([]) ->
+    process_flag(trap_exit, true),
+    ok = config:listen_for_changes(?MODULE, nil),
+    {ok, #st{ctx = plugins_reconcile(#{})}}.
+
+handle_call(reconcile, _From, #st{ctx = Ctx} = St) ->
+    {reply, ok, St#st{ctx = plugins_reconcile(Ctx)}};
+handle_call(context, _From, #st{ctx = Ctx} = St) ->
+    {reply, lists:sort(maps:keys(Ctx)), St};
+handle_call(Msg, _From, #st{} = St) ->
+    {reply, {error, {invalid_call, Msg}}, St}.
+
+handle_cast(reconcile, #st{ctx = Ctx} = St) ->
+    {noreply, St#st{ctx = plugins_reconcile(Ctx)}};
+handle_cast(_Msg, #st{} = St) ->
+    {noreply, St}.
+
+handle_info(restart_config_listener, #st{} = St) ->
+    ok = config:listen_for_changes(?MODULE, nil),
+    {noreply, St};
+handle_info({'EXIT', Pid, Reason}, #st{} = St) ->
+    couch_log:warning("~p: linked process ~p exited: ~p", [?MODULE, Pid, 
Reason]),
+    {noreply, St};
+handle_info(_Msg, #st{} = St) ->
+    {noreply, St}.
+
+terminate(_Reason, #st{ctx = Ctx}) ->
+    _ = lists:foldl(fun plugin_cleanup/2, Ctx, maps:keys(Ctx)),
+    ok.
+
+handle_config_change("replicator", "auth_plugins", _V, _Persist, S) ->
+    ok = gen_server:cast(?MODULE, reconcile),
+    {ok, S};
+handle_config_change(_, _, _, _, S) ->
+    {ok, S}.
+
+handle_config_terminate(_, stop, _) ->
+    ok;
+handle_config_terminate(_, _, _) ->
+    Pid = whereis(?MODULE),
+    erlang:send_after(?RELISTEN_DELAY, Pid, restart_config_listener).
+
+% Private plugin supervisor helper functions
+
+% Start/stop contexts based on `[replicator] auth_plugins` settings.
+plugins_reconcile(Ctx) when is_map(Ctx) ->
+    Wanted = [Mod || Mod <- get_plugin_modules(), exports_sup_initialize(Mod)],
+    Running = maps:keys(Ctx),
+    Ctx1 = lists:foldl(fun plugin_cleanup/2, Ctx, Running -- Wanted),
+    lists:foldl(fun plugin_initialize/2, Ctx1, Wanted -- Running).
+
+plugin_initialize(Mod, Ctx) ->
+    try Mod:sup_initialize() of
+        {ok, PluginCtx} ->
+            couch_log:info("~p: initialized auth plugin ~p", [?MODULE, Mod]),
+            Ctx#{Mod => PluginCtx};
+        ignore ->
+            Ctx;
+        {error, Error} ->
+            LogMsg = "~p: auth plugin ~p sup_initialize failed: ~p",
+            couch_log:error(LogMsg, [?MODULE, Mod, Error]),
+            Ctx
+    catch
+        Tag:Error:Stack ->
+            LogMsg = "~p: auth plugin ~p sup_initialize crashed: ~p:~p ~p",
+            couch_log:error(LogMsg, [?MODULE, Mod, Tag, Error, Stack]),
+            Ctx
+    end.
+
+plugin_cleanup(Mod, Ctx) ->
+    case maps:take(Mod, Ctx) of
+        {PluginCtx, Ctx1} ->
+            % try...catch so if it throw it won't leave a stale entry
+            % A stale entry would block subsequent re-initialization
+            try Mod:sup_cleanup(PluginCtx) of
+                _ -> ok
+            catch
+                Tag:Error:Stack ->
+                    LogMsg = "~p: auth plugin ~p sup_cleanup crashed: ~p:~p 
~p",
+                    couch_log:error(LogMsg, [?MODULE, Mod, Tag, Error, Stack])
+            end,
+            couch_log:info("~p: cleaned up auth plugin ~p", [?MODULE, Mod]),
+            Ctx1;
+        error ->
+            Ctx
+    end.
+
+exports_sup_initialize(Mod) ->
+    % function_exported/3 is false for a module that has not been loaded yet,
+    % and a plugin named in config may not exist at all.
+    case code:ensure_loaded(Mod) of
+        {module, Mod} ->
+            erlang:function_exported(Mod, sup_initialize, 0);
+        {error, Reason} ->
+            LogMsg = "~p: could not load auth plugin module ~p: ~p",
+            couch_log:error(LogMsg, [?MODULE, Mod, Reason]),
+            false
+    end.
 
 -spec get_plugin_modules() -> [atom()].
 get_plugin_modules() ->
diff --git a/src/couch_replicator/src/couch_replicator_sup.erl 
b/src/couch_replicator/src/couch_replicator_sup.erl
index 710919c45..6bd8f669c 100644
--- a/src/couch_replicator/src/couch_replicator_sup.erl
+++ b/src/couch_replicator/src/couch_replicator_sup.erl
@@ -34,7 +34,8 @@ init(_Args) ->
         worker(couch_replicator_connection),
         worker(couch_replicator_rate_limiter),
         worker(couch_replicator_scheduler),
-        worker(couch_replicator_doc_processor)
+        worker(couch_replicator_doc_processor),
+        worker(couch_replicator_auth)
     ],
     SupFlags = #{strategy => rest_for_one, intensity => 10, period => 1},
     {ok, {SupFlags, Children}}.

Reply via email to