|
62 | 62 | keep_alive/2, |
63 | 63 | to_list/2, |
64 | 64 | to_array/2, |
65 | | - parse_mysql_version/2]). |
| 65 | + parse_mysql_version/2, |
| 66 | + get_workers_status/1]). |
66 | 67 |
|
67 | 68 | %% gen_fsm callbacks |
68 | 69 | -export([init/1, handle_event/3, handle_sync_event/4, |
|
106 | 107 | timeout :: pos_integer()}). |
107 | 108 |
|
108 | 109 | -define(STATE_KEY, ejabberd_sql_state). |
| 110 | +-define(STATUS_KEY, ejabberd_sql_status). |
109 | 111 | -define(NESTING_KEY, ejabberd_sql_nesting_level). |
110 | 112 | -define(TOP_LEVEL_TXN, 0). |
111 | 113 | -define(MAX_TRANSACTION_RESTARTS, 10). |
@@ -390,6 +392,29 @@ get_worker(Host) -> |
390 | 392 | get_worker_name(Host, I) -> |
391 | 393 | <<"ejabberd_sql_", Host/binary, $_, (integer_to_binary(I))/binary>>. |
392 | 394 |
|
| 395 | +-spec get_workers_status(binary()) -> #{connected => pos_integer(), disconnected => pos_integer(), overloaded => pos_integer}. |
| 396 | +get_workers_status(Host) -> |
| 397 | + PoolSize = ejabberd_option:sql_pool_size(Host), |
| 398 | + lists:foldl( |
| 399 | + fun(I, #{connected := C, disconnected := D, overloaded := O} = Acc) -> |
| 400 | + maybe |
| 401 | + Pid ?= whereis(binary_to_atom(get_worker_name(Host, I), utf8)), |
| 402 | + true ?= is_pid(Pid), |
| 403 | + {dictionary, Dict} ?= process_info(Pid, dictionary), |
| 404 | + case lists:keyfind(?STATUS_KEY, 1, Dict) of |
| 405 | + {_, connected} -> Acc#{connected => C+1}; |
| 406 | + {_, {overloaded, TS}} -> |
| 407 | + case current_time() - TS > timer:seconds(60) of |
| 408 | + true -> Acc#{connected => C+1}; |
| 409 | + _ -> Acc#{overloaded => O+1} |
| 410 | + end; |
| 411 | + _ -> Acc#{disconnected => D+1} |
| 412 | + end |
| 413 | + else |
| 414 | + _ -> Acc#{disconnected => D+1} |
| 415 | + end |
| 416 | + end, #{connected => 0, disconnected => 0, overloaded => 0}, lists:seq(1, PoolSize)). |
| 417 | + |
393 | 418 | %%%---------------------------------------------------------------------- |
394 | 419 | %%% Callback functions from gen_fsm |
395 | 420 | %%%---------------------------------------------------------------------- |
@@ -437,6 +462,7 @@ connecting(connect, #state{host = Host} = State) -> |
437 | 462 | State1 = State#state{db_ref = Ref, |
438 | 463 | pending_requests = PendingRequests}, |
439 | 464 | State2 = get_db_version(State1), |
| 465 | + put(?STATUS_KEY, connected), |
440 | 466 | {next_state, session_established, State2#state{reconnect_count = 0}} |
441 | 467 | catch _:Reason -> |
442 | 468 | handle_reconnect(Reason, State) |
@@ -548,6 +574,7 @@ handle_reconnect(Reason, #state{host = Host, reconnect_count = RC} = State) -> |
548 | 574 | pgsql -> catch pgsql:terminate(State#state.db_ref); |
549 | 575 | _ -> ok |
550 | 576 | end, |
| 577 | + put(?STATUS_KEY, disconnected), |
551 | 578 | p1_fsm:send_event_after(StartInterval, connect), |
552 | 579 | {next_state, connecting, State#state{reconnect_count = RC + 1, |
553 | 580 | timeout = query_timeout(Host)}}. |
@@ -1039,6 +1066,7 @@ abort_on_driver_error(Reply, From, Timestamp) -> |
1039 | 1066 | -spec report_overload(state()) -> state(). |
1040 | 1067 | report_overload(#state{overload_reported = PrevTime} = State) -> |
1041 | 1068 | CurrTime = current_time(), |
| 1069 | + put(?STATUS_KEY, {overloaded, CurrTime}), |
1042 | 1070 | case PrevTime == undefined orelse (CurrTime - PrevTime) > timer:seconds(30) of |
1043 | 1071 | true -> |
1044 | 1072 | ?ERROR_MSG("SQL connection pool is overloaded, " |
|
0 commit comments