This is an automated email from the ASF dual-hosted git repository. nickva pushed a commit to branch index-info-2 in repository https://gitbox.apache.org/repos/asf/couchdb.git
commit e33ad62407152498a3d2722413616bea864e962d Author: Nick Vatamaniuc <[email protected]> AuthorDate: Wed Sep 2 23:35:11 2026 -0400 Batched per-db _index_info endpoint Implement a new `/{db}/_index_info` endpoint to return the status of all the indexes in the database in a single endpoint. It matches the general pattern of `_dbs_infos` where we return the db info for all databases in a cluster. Just like with `_db_infos`, where it doesn't return anything new just a better and more optimized summary of all dbs, the new endpoint returns all the index status in a single easy to use endpoint. User could, for example, fetch all the design docs, parse them, and with an if/else check call a per-index specific `_info` endpoint such as `/{db}/_design/{ddoc}/_info`, `/{db}/_design/{ddoc}/_nouveau_info/{index}` etc. to get all this information, but that's not very ergonomic and it's slower. One of the main uses of this endpoint is to find out when all the indexes in the db have been built. Each index returns an `updates_pending` object with the pending bounds across all copies of each range. With a minimum (the best copy), maximum (the worst copy) and a count of how many live copies returned a valid result. This information is more accurate than what we might get with the view group `_info` calls. It turns out for view groups `_info` we broke the minimum/preferred/total `updates_pending` stats for view groups. That happened when we switched the collection to work as a ring (the COUCHDB-3053 dedup fix), so then only the first responding copy of each range is aggregated, total = minimum and prefered only counts the responding shards. The new endpoint on the other hand, waits and aggregates stats from all copies of the index. --- src/chttpd/src/chttpd_db.erl | 92 ++++++++ src/chttpd/src/chttpd_httpd_handlers.erl | 1 + src/docs/src/api/database/index.rst | 1 + src/docs/src/api/database/index_info.rst | 224 +++++++++++++++++++ src/docs/src/api/ddoc/common.rst | 4 + src/docs/src/api/ddoc/nouveau.rst | 10 +- src/docs/src/api/ddoc/search.rst | 10 +- src/dreyfus/src/dreyfus_fabric_info.erl | 2 + src/dreyfus/src/dreyfus_index.erl | 34 +++ src/dreyfus/src/dreyfus_rpc.erl | 18 +- src/fabric/src/fabric.erl | 8 + src/fabric/src/fabric_index_info.erl | 372 +++++++++++++++++++++++++++++++ src/fabric/src/fabric_rpc.erl | 129 ++++++++++- src/nouveau/src/nouveau_rpc.erl | 12 +- src/nouveau/src/nouveau_util.erl | 31 +++ 15 files changed, 918 insertions(+), 30 deletions(-) diff --git a/src/chttpd/src/chttpd_db.erl b/src/chttpd/src/chttpd_db.erl index 20f41e04a..ee1b9112c 100644 --- a/src/chttpd/src/chttpd_db.erl +++ b/src/chttpd/src/chttpd_db.erl @@ -26,6 +26,7 @@ handle_changes_req/2, update_doc_result_to_json/1, update_doc_result_to_json/2, handle_design_info_req/3, + handle_index_info_req/2, handle_view_cleanup_req/2, update_doc/4, http_code_from_status/1, @@ -464,6 +465,65 @@ handle_design_info_req(#httpd{method = 'GET'} = Req, Db, #doc{} = DDoc) -> handle_design_info_req(Req, _Db, _DDoc) -> send_method_not_allowed(Req, "GET"). +% $db/_index_info: build status of indexes in a db +handle_index_info_req(#httpd{method = 'GET', path_parts = [_, _]} = Req, Db) -> + Types = parse_types(Req), + DbName = couch_db:name(Db), + {ok, {CopiesExpected, Indexes}} = fabric:get_index_info(DbName, Types), + IndexesEjson = {index_info_to_json(Indexes)}, + Ejson = {[{name, DbName}, {copies_expected, CopiesExpected}, {indexes, IndexesEjson}]}, + send_json(Req, 200, Ejson); +handle_index_info_req(#httpd{method = 'GET'} = Req, _Db) -> + chttpd:send_error(Req, not_found); +handle_index_info_req(Req, _Db) -> + send_method_not_allowed(Req, "GET"). + +% _index_info types: one or more index types. ex: ?type=nouveau,search +parse_types(Req) -> + TypeArgs = [V || {K, V} <- chttpd:qs(Req), K == "type"], + case TypeArgs of + [] -> + [view, search, nouveau]; + [_ | _] -> + Tokens = lists:flatmap(fun(V) -> string:tokens(V, ",") end, TypeArgs), + case lists:usort(lists:map(fun parse_type/1, Tokens)) of + [] -> throw({query_parse_error, <<"`type` must not be empty">>}); + Types -> [T || T <- [view, search, nouveau], lists:member(T, Types)] + end + end. + +parse_type("view") -> + view; +parse_type("search") -> + search; +parse_type("nouveau") -> + nouveau; +parse_type(Other) -> + Msg = io_lib:format("Invalid index type: ~s. Must be view, search or nouveau", [Other]), + throw({query_parse_error, ?l2b(Msg)}). + +index_info_to_json(Indexes) -> + [{DDocId, index_info(Res)} || {DDocId, Res} <- Indexes]. + +index_info({error, Error}) -> + errobj(Error); +index_info({Sections}) -> + {[{Sect, index_section_json(Sect, Val)} || {Sect, Val} <- Sections]}. + +index_section_json(view_index, Res) -> + res_to_json(Res); +index_section_json(_Named, {Idxs}) -> + {[{Name, res_to_json(Res)} || {Name, Res} <- Idxs]}. + +res_to_json({error, Error}) -> + errobj(Error); +res_to_json(Info) -> + Info. + +errobj(Error) -> + {_Code, ErrorStr, ReasonStr} = chttpd:error_info(Error), + {[{error, ErrorStr}, {reason, ReasonStr}]}. + create_db_req(#httpd{} = Req, DbName) -> couch_httpd:verify_is_server_admin(Req), ShardsOpt = parse_shards_opt(Req), @@ -2660,6 +2720,38 @@ monitor_attachments_test_() -> ?_assertEqual([], monitor_attachments(Atts)) end}. +index_info_to_json_test_() -> + % {error, Reason} turn into {"error":..., "reason":...} + Unavailable = {service_unavailable, <<"Search is not available">>}, + Indexes = [ + {<<"_design/err">>, {error, not_found}}, + {<<"_design/ok">>, + {[ + {view_index, {[{update_seq, 1}]}}, + {search_indexes, + {[{<<"i">>, {[{doc_count, 2}]}}, {<<"j">>, {error, Unavailable}}]}}, + {nouveau_indexes, {[{<<"n">>, #{num_docs => 3}}]}} + ]}} + ], + Expected = [ + {<<"_design/err">>, {[{error, <<"not_found">>}, {reason, <<"missing">>}]}}, + {<<"_design/ok">>, + {[ + {view_index, {[{update_seq, 1}]}}, + {search_indexes, + {[ + {<<"i">>, {[{doc_count, 2}]}}, + {<<"j">>, + {[ + {error, <<"service unavailable">>}, + {reason, <<"Search is not available">>} + ]}} + ]}}, + {nouveau_indexes, {[{<<"n">>, #{num_docs => 3}}]}} + ]}} + ], + ?_assertEqual(Expected, index_info_to_json(Indexes)). + parse_partitioned_opt_test_() -> { foreach, diff --git a/src/chttpd/src/chttpd_httpd_handlers.erl b/src/chttpd/src/chttpd_httpd_handlers.erl index 3e499b72d..4251cbfbc 100644 --- a/src/chttpd/src/chttpd_httpd_handlers.erl +++ b/src/chttpd/src/chttpd_httpd_handlers.erl @@ -30,6 +30,7 @@ url_handler(<<"_up">>) -> fun chttpd_misc:handle_up_req/1; url_handler(_) -> no_match. db_handler(<<"_view_cleanup">>) -> fun chttpd_db:handle_view_cleanup_req/2; +db_handler(<<"_index_info">>) -> fun chttpd_db:handle_index_info_req/2; db_handler(<<"_compact">>) -> fun chttpd_db:handle_compact_req/2; db_handler(<<"_design">>) -> fun chttpd_db:handle_design_req/2; db_handler(<<"_partition">>) -> fun chttpd_db:handle_partition_req/2; diff --git a/src/docs/src/api/database/index.rst b/src/docs/src/api/database/index.rst index 914c0d6ff..8b5f23c1c 100644 --- a/src/docs/src/api/database/index.rst +++ b/src/docs/src/api/database/index.rst @@ -44,5 +44,6 @@ Where ``{db}`` is the name of any database. changes compact cleanup + index_info security misc diff --git a/src/docs/src/api/database/index_info.rst b/src/docs/src/api/database/index_info.rst new file mode 100644 index 000000000..a862f1b73 --- /dev/null +++ b/src/docs/src/api/database/index_info.rst @@ -0,0 +1,224 @@ +.. Licensed under the Apache License, Version 2.0 (the "License"); you may not +.. use this file except in compliance with the License. You may obtain a copy of +.. the License at +.. +.. http://www.apache.org/licenses/LICENSE-2.0 +.. +.. Unless required by applicable law or agreed to in writing, software +.. distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +.. WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +.. License for the specific language governing permissions and limitations under +.. the License. + +.. _api/db/index_info: + +===================== +``/{db}/_index_info`` +===================== + +.. versionadded:: 3.6 + +.. http:get:: /{db}/_index_info + :synopsis: Returns status all the indexes in the database + + Index info objects have the same fields as the ones returned by + :get:`/{db}/_design/{ddoc}/_info`, + :get:`/{db}/_design/{ddoc}/_search_info/{index}` and + :get:`/{db}/_design/{ddoc}/_nouveau_info/{index}`, with an + ``updates_pending`` object with the index build pending bounds across all + copies of the database. See :ref:`api/db/index_info/structure`. + + This endpoint may be used when one needs to know whether a database's + indexes are fully built, for example before switching traffic to a replica. + + :param db: Database name + :query string type: Filter response to the given index types. Can be one of + these values: ``view``, ``search`` and ``nouveau``, comma separated. + *Optional*, default: all types. + :<header Accept: - :mimetype:`application/json` + - :mimetype:`text/plain` + :>header Content-Type: - :mimetype:`application/json` + - :mimetype:`text/plain; charset=utf-8` + :>json string name: Database name + :>json number copies_expected: The number of shard copies of the database + (Q*N) Compare with the ``copies`` count of each ``updates_pending`` + object. + :>json object indexes: Index information keyed by design doc ID. + See :ref:`api/db/index_info/structure`. + :code 200: Request completed successfully + :code 400: Invalid database name or ``type`` value + :code 401: Unauthorized request to a protected API + :code 403: Insufficient permissions / :ref:`Too many requests with invalid credentials<error/403>` + :code 404: Database doesn't exist + + **Request**: + + .. code-block:: http + + GET /recipes/_index_info HTTP/1.1 + Accept: application/json + Host: localhost:5984 + + **Response**: + + .. code-block:: http + + HTTP/1.1 200 OK + Cache-Control: must-revalidate + Content-Type: application/json + Date: Mon, 01 Sep 2025 15:42:11 GMT + Server: CouchDB (Erlang/OTP) + + { + "name": "recipes", + "copies_expected": 6, + "indexes": { + "_design/cookbook": { + "view_index": { + "collator_versions": [ + "153.136" + ], + "compact_running": false, + "language": "javascript", + "purge_seq": 0, + "signature": "a6d97b0199e54a1eb56e4becb1322587", + "sizes": { + "active": 1825, + "external": 1355, + "file": 16750 + }, + "update_seq": 54, + "updater_running": false, + "waiting_clients": 0, + "waiting_commit": false, + "updates_pending": { + "minimum": 0, + "maximum": 0, + "copies": 6 + } + }, + "search_indexes": { + "ingredients": { + "committed_seq": 54, + "disk_size": 3960, + "doc_count": 50, + "doc_del_count": 0, + "pending_seq": 54, + "signature": "0b4ba635d5eb4fcbb2f6c9c2247460ec", + "updates_pending": { + "minimum": 0, + "maximum": 0, + "copies": 6 + } + } + }, + "nouveau_indexes": { + "ingredients": { + "disk_size": 6324, + "num_docs": 50, + "purge_seq": 0, + "signature": "ea87fe8f9517403691850f51d0a1ce3e5afaf89347204dd79430252a0591e503", + "update_seq": 54, + "updates_pending": { + "minimum": 0, + "maximum": 0, + "copies": 6 + } + } + } + }, + "_design/8c2a4caf8ea1b581ac43a062fd43a876dee1382d": { + "view_index": { + "collator_versions": [ + "153.136" + ], + "compact_running": false, + "language": "query", + "purge_seq": 0, + "signature": "b77547252cb8b19ff12831973b576c0f", + "sizes": { + "active": 0, + "external": 10, + "file": 102 + }, + "update_seq": 0, + "updater_running": false, + "waiting_clients": 0, + "waiting_commit": false, + "updates_pending": { + "minimum": 54, + "maximum": 54, + "copies": 6 + } + } + } + } + } + + In this example the ``cookbook`` design document has a view group, a + search index and a nouveau index, all fully built on all six copies. The + second design document is a Mango index which has not been built on + any copy yet: every copy is 54 updates behind. + +.. _api/db/index_info/structure: + +Index Information +================= + +The ``indexes`` object of the :get:`/{db}/_index_info` response has one entry +per design document with at least one index, keyed by design doc ID. Each entry +may contain: + +* **view_index** (*object*): The design doc's view group, with the fields + described in :ref:`api/ddoc/view_index_info` and ``updates_pending``. A design + document has one view group for all of its views. + +* **search_indexes** (*object*): One entry per search index of the design + document, keyed by index name, with the fields of the ``search_index`` object + returned by :get:`/{db}/_design/{ddoc}/_search_info/{index}` plus + ``updates_pending``. +* **nouveau_indexes** (*object*): One entry per nouveau index of the design + document, keyed by index name, with the fields of the ``search_index`` object + returned by :get:`/{db}/_design/{ddoc}/_nouveau_info/{index}` plus + ``updates_pending``. + +Pending stats +-------------- + +Every index object contains an ``updates_pending`` object which shows the +number of pending changes for that index. + +* **minimum** (*number*): Backlog of the most up to date copy of each + range. ``0`` means at least one fully built copy of the index exists for + every range. +* **maximum** (*number*): Backlog of the least up to date copy of each range. + ``0`` means the index is fully built on all returned shard copies. +* **copies** (*number*): The number of shard copies with returned a response. + Copies which are unreachable, in maintenance mode, or which did not answer + before the request timed out are skipped. When this is much smaller than + ``copies_expected`` the bounds are partial and ``maximum`` may underestimate + the true pending backlog. + +Errors +------ + +If an index cannot be inspected it will return an error for that index only. +For example is Closeau is not available search indexes will return: + +.. code-block:: javascript + + "search_indexes": { + "ingredients": { + "error": "service unavailable", + "reason": "Search is not available" + } + } + +Invalid ``type`` values return a ``400 Bad Request``: + +.. code-block:: javascript + + { + "error": "query_parse_error", + "reason": "Invalid index type: foo. Must be view, search or nouveau" + } diff --git a/src/docs/src/api/ddoc/common.rst b/src/docs/src/api/ddoc/common.rst index 961af135b..d8193bcef 100644 --- a/src/docs/src/api/ddoc/common.rst +++ b/src/docs/src/api/ddoc/common.rst @@ -148,6 +148,10 @@ index, index size and current status of the design document and associated index information. + To obtain this information for every design document of a database, + including search and nouveau indexes, in a single request see + :get:`/{db}/_index_info`. + :param db: Database name :param ddoc: Design document name :<header Accept: - :mimetype:`application/json` diff --git a/src/docs/src/api/ddoc/nouveau.rst b/src/docs/src/api/ddoc/nouveau.rst index fa9935e97..197722d14 100644 --- a/src/docs/src/api/ddoc/nouveau.rst +++ b/src/docs/src/api/ddoc/nouveau.rst @@ -143,6 +143,14 @@ "search_index": { "num_docs": 1000, "update_seq": 5000, - "disk_size": 1048576 + "disk_size": 1048576, + "pending_updates": 100, + "signature": "ef30790cdaf74ad78daebee9597a17393d875db085184cf688aa0e8abcb972f2" } } + +.. versionadded:: 3.6 + ``pending_updates`` is the number of database updates the index still + has to process, summed over shard ranges. ``0`` means the index is up to + date with the database. To obtain it for every index of a database in one + request see :get:`/{db}/_index_info`. diff --git a/src/docs/src/api/ddoc/search.rst b/src/docs/src/api/ddoc/search.rst index 427a9b507..869e74600 100644 --- a/src/docs/src/api/ddoc/search.rst +++ b/src/docs/src/api/ddoc/search.rst @@ -169,6 +169,14 @@ "doc_del_count": 129180, "doc_count": 1066173, "disk_size": 728305827, - "committed_seq": 7125496 + "committed_seq": 7125496, + "pending_updates": 0, + "signature": "13083011f9554446a4ac093927e75d07" } } + +.. versionadded:: 3.6 + ``pending_updates`` is the number of database updates the index still + has to process, summed over shard ranges. ``0`` means the index is up to + date with the database. To obtain it for every index of a database in one + request see :get:`/{db}/_index_info`. diff --git a/src/dreyfus/src/dreyfus_fabric_info.erl b/src/dreyfus/src/dreyfus_fabric_info.erl index 2925768b1..d7eae223a 100644 --- a/src/dreyfus/src/dreyfus_fabric_info.erl +++ b/src/dreyfus/src/dreyfus_fabric_info.erl @@ -100,6 +100,8 @@ merge_results(Info) -> [{committed_seq, lists:sum(X)} | Acc]; (pending_seq, X, Acc) -> [{pending_seq, lists:sum(X)} | Acc]; + (pending_updates, X, Acc) -> + [{pending_updates, lists:sum(X)} | Acc]; (signature, [X | _], Acc) -> [{signature, X} | Acc]; (_, _, Acc) -> diff --git a/src/dreyfus/src/dreyfus_index.erl b/src/dreyfus/src/dreyfus_index.erl index d75405001..29afaedb4 100644 --- a/src/dreyfus/src/dreyfus_index.erl +++ b/src/dreyfus/src/dreyfus_index.erl @@ -26,6 +26,7 @@ await/2, search/2, info/1, + info/3, group1/2, group2/2, design_doc_to_indexes/2 @@ -79,6 +80,39 @@ info(Pid0) -> MFA = {?MODULE, info_int, [Pid]}, dreyfus_util:time([index, info], MFA). +% Get shard local index info based on ddoc and index name. This can be used by +% _index_info worker when getting index info batches. +info(DbName, DDoc, IndexName) -> + case design_doc_to_index(DbName, DDoc, IndexName) of + {ok, Index} -> + case dreyfus_index_manager:get_index(DbName, Index) of + {ok, Pid} -> + case info(Pid) of + {ok, Fields} -> + Info = [ + {signature, Index#index.sig}, + {pending_updates, pending_updates(DbName, Fields)} + | Fields + ], + {ok, Info}; + Else -> + Else + end; + Error -> + Error + end; + Error -> + Error + end. + +% Similar idea as view pending_updates from couch_index. Return a difference +% betweenn committed db seq and index +pending_updates(DbName, Fields) -> + GetCommSeq = fun(Db) -> couch_db:get_committed_update_seq(Db) end, + CommittedSeq = couch_util:with_db(DbName, GetCommSeq), + IndexSeq = couch_util:get_value(pending_seq, Fields, 0), + max(CommittedSeq - IndexSeq, 0). + %% We either have a dreyfus_index gen_server pid or the remote %% clouseau pid. to_index_pid(Pid) -> diff --git a/src/dreyfus/src/dreyfus_rpc.erl b/src/dreyfus/src/dreyfus_rpc.erl index 5fbe194dc..b00af49e6 100644 --- a/src/dreyfus/src/dreyfus_rpc.erl +++ b/src/dreyfus/src/dreyfus_rpc.erl @@ -80,23 +80,7 @@ info(DbName, DDoc, IndexName) -> info_int(DbName, DDoc, IndexName) -> erlang:put(io_priority, {search, DbName}), check_interactive_mode(), - case dreyfus_index:design_doc_to_index(DbName, DDoc, IndexName) of - {ok, Index} -> - case dreyfus_index_manager:get_index(DbName, Index) of - {ok, Pid} -> - case dreyfus_index:info(Pid) of - {ok, Fields} -> - Info = [{signature, Index#index.sig} | Fields], - rexi:reply({ok, Info}); - Else -> - rexi:reply(Else) - end; - Error -> - rexi:reply(Error) - end; - Error -> - rexi:reply(Error) - end. + rexi:reply(dreyfus_index:info(DbName, DDoc, IndexName)). disk_size(DbName, DDoc, IndexName) -> erlang:put(io_priority, {search, DbName}), diff --git a/src/fabric/src/fabric.erl b/src/fabric/src/fabric.erl index ab575f135..5099ab519 100644 --- a/src/fabric/src/fabric.erl +++ b/src/fabric/src/fabric.erl @@ -59,6 +59,7 @@ changes/4, query_view/3, query_view/4, query_view/6, query_view/7, get_view_group_info/2, + get_index_info/2, end_changes/0 ]). @@ -595,6 +596,13 @@ query_view(Db, Options, DDoc, ViewName, Callback, Acc0, QueryArgs0) -> get_view_group_info(DbName, DesignId) -> fabric_group_info:go(dbname(DbName), design_doc(DesignId)). +%% @doc get index info of all indexes: view, search and nouveau etc +%% +-spec get_index_info(dbname(), [view | search | nouveau]) -> + {ok, {non_neg_integer(), [{binary(), any()}]}}. +get_index_info(DbName, Types) -> + fabric_index_info:go(dbname(DbName), Types). + -spec end_changes() -> ok. end_changes() -> fabric_view_changes:increment_changes_epoch(). diff --git a/src/fabric/src/fabric_index_info.erl b/src/fabric/src/fabric_index_info.erl new file mode 100644 index 000000000..433075ce7 --- /dev/null +++ b/src/fabric/src/fabric_index_info.erl @@ -0,0 +1,372 @@ +% Licensed under the Apache License, Version 2.0 (the "License"); you may not +% use this file except in compliance with the License. You may obtain a copy of +% the License at +% +% http://www.apache.org/licenses/LICENSE-2.0 +% +% Unless required by applicable law or agreed to in writing, software +% distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +% WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +% License for the specific language governing permissions and limitations under +% the License. + +% Get status of all the indexes in db. Info gathering is optimized by sending +% the ddocs to the works and they locally gather index info in parallel. +% +% To get an idea about how many shards are built we query every live copy and +% return a summary of expected shards vs actual shards included in the results. +% The result returns the expected number of shards (workers) and each index +% summary will contain the number of copies which actually replied. This gives +% us an idea how many indexes and how far along they are before they are built. +% +% For each index besides their sizes and other info we return: +% +% minimum : best copy of each range. When 0 it means at least one complete +% copy of every range exists. +% +% maximum : the worst copy of each range. When 0 it means every copy of every +% live range is fully bult. +% +% copies : live shard copies which returned results. If this value is too low +% compared to the expected copies the maximum/minimum values may not be +% trusted much. +% +% For other index info fields use the lowest values per range so they are +% more or less stable from call to call. + +-module(fabric_index_info). + +-export([go/2]). + +-include_lib("mem3/include/mem3.hrl"). + +go(DbName, Types) -> + {ok, DDocs} = fabric:design_docs(DbName), + {CopiesExpected, Responses} = index_info(DbName, DDocs, Types), + {ok, {CopiesExpected, merge_responses(Responses)}}. + +index_info(DbName, DDocs, Types) -> + Shards = mem3:shards(DbName), + Workers = fabric_util:submit_jobs(Shards, fabric_rpc, index_info, [DDocs, Types]), + RexiMon = fabric_util:create_monitors(Workers), + Acc0 = {fabric_dict:init(Workers, nil), []}, + try fabric_util:recv(Workers, #shard.ref, fun handle_message/3, Acc0) of + {ok, Responses} -> + {length(Workers), tag_responses(Responses)}; + {timeout, {WorkersDict, Responses}} -> + DefunctWorkers = fabric_util:remove_done_workers(WorkersDict, nil), + fabric_util:log_timeout(DefunctWorkers, "index_info"), + fabric_util:cleanup(DefunctWorkers), + {length(Workers), tag_responses(Responses)}; + {error, Error} -> + fabric_util:cleanup(Workers), + throw(Error) + after + rexi_monitor:stop(RexiMon) + end. + +tag_responses(Responses) -> + [{W#shard.range, W#shard.node, PerDDoc} || {W, PerDDoc} <- Responses]. + +% Failed / unreachables dropped. Otherwise responses are merged +% +handle_message({ok, Result}, Worker, {Counters, Acc}) -> + Counters1 = fabric_dict:erase(Worker, Counters), + maybe_stop(Counters1, [{Worker, Result} | Acc]); +handle_message({rexi_DOWN, _, {_, NRef}, _}, _Worker, {Counters, Acc}) -> + Counters1 = fabric_dict:filter(fun(#shard{node = N}, _) -> N =/= NRef end, Counters), + maybe_stop(Counters1, Acc); +handle_message(_Error, Worker, {Counters, Acc}) -> + Counters1 = fabric_dict:erase(Worker, Counters), + maybe_stop(Counters1, Acc). + +maybe_stop(Counters, Acc) -> + case fabric_dict:size(Counters) of + 0 -> {stop, Acc}; + _ -> {ok, {Counters, Acc}} + end. + +% Each response is {Range, Node, [{DDocId, Sections}]} +% Sections: {error, Err} | [ +% {view_index, Res}, +% {search_indexes, [{Name, Res}]}, +% {nouveau_indexes, [{Name, Res}]} +% ] +% Res: {ok, Info} | {error, Err} +% +% Note: view ddoc is a group so there only one per ddoc +% +merge_responses(Responses) -> + ByDDoc = lists:foldl( + fun({Range, Node, PerDDoc}, Acc0) -> + lists:foldl( + fun({DDocId, Sect}, Acc) -> + orddict:append(DDocId, {Range, Node, Sect}, Acc) + end, + Acc0, + PerDDoc + ) + end, + orddict:new(), + Responses + ), + [{DDocId, merge_ddoc(Entries)} || {DDocId, Entries} <- ByDDoc]. + +merge_ddoc(Entries) -> + case [Error || {_, _, {error, Error}} <- Entries] of + [Error | _] -> + {error, Error}; + [] -> + BySection = lists:foldl( + fun({Range, Node, Sect}, Acc0) -> + lists:foldl( + fun({K, V}, Acc) -> orddict:append(K, {Range, Node, V}, Acc) end, + Acc0, + Sect + ) + end, + orddict:new(), + Entries + ), + {[merge_section(K, Vs) || {K, Vs} <- BySection]} + end. + +merge_section(view_index, Results) -> + {view_index, merge_leaf(view, Results)}; +merge_section(search_indexes, PerCopyIdxLists) -> + {search_indexes, {merge_named(search, PerCopyIdxLists)}}; +merge_section(nouveau_indexes, PerCopyIdxLists) -> + {nouveau_indexes, {merge_named(nouveau, PerCopyIdxLists)}}. + +merge_named(Type, PerCopyIdxLists) -> + ByName = lists:foldl( + fun({Range, Node, IdxList}, Acc0) -> + lists:foldl( + fun({Name, Res}, Acc) -> orddict:append(Name, {Range, Node, Res}, Acc) end, + Acc0, + IdxList + ) + end, + orddict:new(), + PerCopyIdxLists + ), + [{Name, merge_leaf(Type, Results)} || {Name, Results} <- ByName]. + +% Results: {Range, Node, {ok, Info} | {error, Err}} per copy. Errors "win" +merge_leaf(Type, Results) -> + case [Error || {_, _, {error, Error}} <- Results] of + [Error | _] -> {error, Error}; + [] -> merge_oks(Type, Results) + end. + +% Non pending updates are merged by lowest copy. We aggregate pending at the +% end since we're doing a min/max over those. +merge_oks(Type, Results) -> + ByRange = lists:foldl( + fun({Range, Node, {ok, Info}}, Acc) -> orddict:append(Range, {Node, Info}, Acc) end, + orddict:new(), + Results + ), + RangeCopies = [Copies || {_Range, Copies} <- ByRange], + Reps = [rem_pending(Type, first(Copies)) || Copies <- RangeCopies], + PerRange = [[pending(Type, Info) || {_Node, Info} <- Copies] || Copies <- RangeCopies], + add_pending(Type, merge_infos(Type, Reps), agg_pending(PerRange)). + +% For some stable-ish order results don't flip flop with each call +first(Copies) -> + [{_Node, Info} | _] = lists:sort(Copies), + Info. + +% Skip copies which did not report pending +% Ranges with no reporting copies contribute 0 +% Return undefined when no copy of any range sent reports +agg_pending(PerRange) -> + case lists:foldl(fun pending_range/2, undefined, PerRange) of + undefined -> undefined; + {Min, Max, Copies} -> [{minimum, Min}, {maximum, Max}, {copies, Copies}] + end. + +pending_range(Vals, Acc) -> + case lists:sort([P || P <- Vals, is_integer(P)]) of + [] -> + Acc; + [_ | _] = Vs -> + {Min, Max, Copies} = + case Acc of + undefined -> {0, 0, 0}; + _ -> Acc + end, + {Min + hd(Vs), Max + lists:last(Vs), Copies + length(Vs)} + end. + +merge_infos(view, Infos) -> + merge_view_infos(Infos); +merge_infos(search, Infos) -> + merge_search_infos(Infos); +merge_infos(nouveau, Infos) -> + merge_nouveau_infos(Infos). + +pending(nouveau, #{} = Info) -> + maps:get(pending_updates, Info, undefined); +pending(_Type, Info) -> + couch_util:get_value(pending_updates, Info). + +rem_pending(nouveau, #{} = Info) -> + maps:remove(pending_updates, Info); +rem_pending(_Type, Info) -> + lists:keydelete(pending_updates, 1, Info). + +add_pending(_Type, Merged, undefined) -> + Merged; +add_pending(nouveau, #{} = Merged, Bounds) -> + Merged#{updates_pending => maps:from_list(Bounds)}; +add_pending(_Type, {Props}, Bounds) -> + {lists:keystore(updates_pending, 1, Props, {updates_pending, {Bounds}})}. + +% Copied from fabric_group_info:merge_results/1 mostly. Except we remove the pending updates +% then do our own (hopefully better) calculation over them. +merge_view_infos(Infos) -> + Dict = to_orddict(lists:append(Infos)), + Merged = orddict:fold( + fun + (signature, [X | _], Acc) -> + [{signature, X} | Acc]; + (language, [X | _], Acc) -> + [{language, X} | Acc]; + (sizes, X, Acc) -> + [{sizes, {merge_obj(X)}} | Acc]; + (compact_running, X, Acc) -> + [{compact_running, lists:member(true, X)} | Acc]; + (updater_running, X, Acc) -> + [{updater_running, lists:member(true, X)} | Acc]; + (waiting_commit, X, Acc) -> + [{waiting_commit, lists:member(true, X)} | Acc]; + (waiting_clients, X, Acc) -> + [{waiting_clients, lists:sum(X)} | Acc]; + (update_seq, X, Acc) -> + [{update_seq, lists:sum(X)} | Acc]; + (purge_seq, X, Acc) -> + [{purge_seq, lists:sum(X)} | Acc]; + (collator_versions, X, Acc) -> + Vs = lists:usort(lists:flatmap(fun(V) -> V end, X)), + [{collator_versions, Vs} | Acc]; + (_, _, Acc) -> + Acc + end, + [], + Dict + ), + {Merged}. + +% Copied from dreyfus_fabric_info:merge_results/1. We remove pending updates +% and add them later. +merge_search_infos(Infos) -> + Dict = to_orddict(lists:append(Infos)), + Merged = orddict:fold( + fun + (signature, [X | _], Acc) -> [{signature, X} | Acc]; + (disk_size, X, Acc) -> [{disk_size, lists:sum(X)} | Acc]; + (doc_count, X, Acc) -> [{doc_count, lists:sum(X)} | Acc]; + (doc_del_count, X, Acc) -> [{doc_del_count, lists:sum(X)} | Acc]; + (committed_seq, X, Acc) -> [{committed_seq, lists:sum(X)} | Acc]; + (pending_seq, X, Acc) -> [{pending_seq, lists:sum(X)} | Acc]; + (_, _, Acc) -> Acc + end, + [], + Dict + ), + {Merged}. + +merge_nouveau_infos(Maps) -> + lists:foldl(fun(M, Acc) -> maps:merge_with(fun merge_nouveau_val/3, M, Acc) end, #{}, Maps). + +merge_nouveau_val(signature, Val, Val) -> + % Can't sum signatures, but we can sum everything else + Val; +merge_nouveau_val(_Key, Val1, Val2) -> + Val1 + Val2. + +to_orddict(KVs) -> + lists:foldl(fun({K, V}, D) -> orddict:append(K, V, D) end, orddict:new(), KVs). + +% Merge list of {[{K, V}]} objects where V is a number +merge_obj(Objects) -> + Dict = lists:foldl( + fun({Props}, D) -> + lists:foldl(fun({K, V}, D0) -> orddict:append(K, V, D0) end, D, Props) + end, + orddict:new(), + Objects + ), + orddict:fold( + fun(Key, X, Acc) -> + [{Key, lists:sum(X)} | Acc] + end, + [], + Dict + ). + +-ifdef(TEST). +-include_lib("couch/include/couch_eunit.hrl"). + +aggregate_pending_test() -> + ?assertEqual(undefined, agg_pending([])), + ?assertEqual(undefined, agg_pending([[]])), + ?assertEqual(undefined, agg_pending([[undefined]])), + ?assertEqual([{minimum, 5}, {maximum, 5}, {copies, 1}], agg_pending([[5]])), + ?assertEqual( + [{minimum, 3}, {maximum, 11}, {copies, 5}], + agg_pending([[4, 0, 8], [3, 3]]) + ), + ?assertEqual( + [{minimum, 2}, {maximum, 2}, {copies, 1}], + agg_pending([[undefined, 2], [undefined]]) + ). + +viewres(Pending, Seq) -> + {ok, [{signature, <<"sig">>}, {pending_updates, Pending}, {update_seq, Seq}]}. + +merge_responses_test() -> + % Some field are added while for pending we compute our min/max stats + Responses = [ + {[0, 10], n1, [{<<"_design/d">>, [{view_index, viewres(0, 100)}]}]}, + {[0, 10], n2, [{<<"_design/d">>, [{view_index, viewres(6, 94)}]}]}, + {[11, 20], n2, [{<<"_design/d">>, [{view_index, viewres(2, 98)}]}]}, + {[11, 20], n3, [{<<"_design/d">>, [{view_index, viewres(4, 96)}]}]} + ], + [{<<"_design/d">>, {[{view_index, {Props}}]}}] = merge_responses(Responses), + ?assertEqual(198, couch_util:get_value(update_seq, Props)), + ?assertEqual(undefined, couch_util:get_value(pending_updates, Props)), + ?assertEqual( + {[{minimum, 2}, {maximum, 10}, {copies, 4}]}, + couch_util:get_value(updates_pending, Props) + ). + +merge_responses_error_wins_test() -> + Responses = [ + {[0, 10], n1, [{<<"_design/d">>, [{view_index, viewres(0, 100)}]}]}, + {[0, 10], n2, [{<<"_design/d">>, [{view_index, {error, not_found}}]}]} + ], + ?assertEqual( + [{<<"_design/d">>, {[{view_index, {error, not_found}}]}}], + merge_responses(Responses) + ). + +merge_responses_nouveau_test() -> + % Nouveau uses maps for info so add a separate test for it + NRes = fun(Pending, Seq) -> + {ok, #{signature => <<"s">>, pending_updates => Pending, <<"update_seq">> => Seq}} + end, + Responses = [ + {[0, 10], n1, [{<<"_design/d">>, [{nouveau_indexes, [{<<"n">>, NRes(5, 10)}]}]}]}, + {[0, 10], n2, [{<<"_design/d">>, [{nouveau_indexes, [{<<"n">>, NRes(0, 15)}]}]}]} + ], + [{<<"_design/d">>, {[{nouveau_indexes, {[{<<"n">>, Merged}]}}]}}] = + merge_responses(Responses), + ?assertEqual(10, maps:get(<<"update_seq">>, Merged)), + ?assertNot(maps:is_key(pending_updates, Merged)), + ?assertEqual( + #{minimum => 0, maximum => 5, copies => 2}, maps:get(updates_pending, Merged) + ). + +-endif. diff --git a/src/fabric/src/fabric_rpc.erl b/src/fabric/src/fabric_rpc.erl index fd06b391b..0553f6190 100644 --- a/src/fabric/src/fabric_rpc.erl +++ b/src/fabric/src/fabric_rpc.erl @@ -26,7 +26,7 @@ get_missing_revs/2, get_missing_revs/3, update_docs/3 ]). --export([all_docs/3, changes/3, map_view/4, reduce_view/4, group_info/2]). +-export([all_docs/3, changes/3, map_view/4, reduce_view/4, group_info/2, index_info/3]). -export([ create_db/1, create_db/2, delete_db/1, @@ -66,6 +66,9 @@ -include_lib("couch/include/couch_db.hrl"). -include_lib("couch_mrview/include/couch_mrview.hrl"). +% How many design docs an index_info/3 worker processes concurrently +-define(MAX_CONCURRENCY, 8). + %% rpc endpoints %% call to with_db will supply your M:F with a Db instance %% and then remaining args @@ -337,6 +340,130 @@ purge_docs(DbName, UUIdsIdsRevs, Options) -> group_info(DbName, DDocId) -> group_info(DbName, DDocId, []). +% Shard local worker _index_info. Get info for all ddocs on this shard copy. +index_info(DbName, DDocs, Types) -> + set_io_priority(DbName, []), + Avail = {search_available(), nouveau_enabled()}, + Infos = pmap(fun(DDoc) -> ddoc_info(DbName, DDoc, Types, Avail) end, DDocs), + rexi:reply({ok, lists:append(Infos)}). + +% Return [{DDocId, Sections}] or [] for empty ddocs +ddoc_info(ShardName, {Props} = DDocEJson, Types, Avail) -> + DDocId = couch_util:get_value(<<"_id">>, Props), + try + #doc{body = {Body}} = DDoc = couch_doc:from_json_obj(DDocEJson), + Sections = lists:append([section(T, ShardName, DDoc, Body, Avail) || T <- Types]), + case Sections of + [] -> []; + [_ | _] -> [{DDocId, Sections}] + end + catch + _Tag:Error -> + [{DDocId, {error, Error}}] + end. + +section(view, ShardName, DDoc, Body, _Avail) -> + case couch_util:get_value(<<"views">>, Body) of + {[_ | _]} -> [{view_index, view_info(ShardName, DDoc)}]; + _ -> [] + end; +section(search, ShardName, DDoc, Body, {SearchAvailable, _}) -> + case couch_util:get_value(<<"indexes">>, Body) of + {[_ | _] = Indexes} -> + Infos = [ + {Name, search_info(SearchAvailable, ShardName, DDoc, Name)} + || {Name, _} <- Indexes + ], + [{search_indexes, Infos}]; + _ -> + [] + end; +section(nouveau, ShardName, DDoc, Body, {_, NouveauEnabled}) -> + case couch_util:get_value(<<"nouveau">>, Body) of + {[_ | _] = Indexes} -> + Infos = [ + {Name, nouveau_info(NouveauEnabled, ShardName, DDoc, Name)} + || {Name, _} <- Indexes + ], + [{nouveau_indexes, Infos}]; + _ -> + [] + end. + +% Same node local call as group_info/3 above. An index which cannot be +% opened is reported as {error, Reason} instead of crashing the worker. +view_info(ShardName, DDoc) -> + try + norm_error(couch_mrview:get_info(ShardName, DDoc)) + catch + error:{badmatch, Error} -> norm_error(Error); + _Tag:Error -> {error, Error} + end. + +search_info(false, _ShardName, _DDoc, _IndexName) -> + {error, {service_unavailable, <<"Search is not available">>}}; +search_info(true, ShardName, DDoc, IndexName) -> + try + norm_error(dreyfus_index:info(ShardName, DDoc, IndexName)) + catch + _Tag:Error -> {error, Error} + end. + +nouveau_info(false, _ShardName, _DDoc, _IndexName) -> + {error, {service_unavailable, <<"nouveau is not enabled">>}}; +nouveau_info(true, ShardName, DDoc, IndexName) -> + try + norm_error(nouveau_util:index_info(ShardName, DDoc, IndexName)) + catch + _Tag:Error -> {error, Error} + end. + +norm_error({ok, Info}) -> + {ok, Info}; +norm_error({error, Error}) -> + {error, Error}; +norm_error(Else) -> + {error, Else}. + +search_available() -> + try + dreyfus:available() + catch + _:_ -> false + end. + +nouveau_enabled() -> + try + nouveau:enabled() + catch + _:_ -> false + end. + +% Run fun with bounded concurrency and isolated in separate processes +% +pmap(Fun, Items) -> + lists:append([pmap_chunk(Fun, Chunk) || Chunk <- chunk(Items, ?MAX_CONCURRENCY)]). + +pmap_chunk(Fun, Items) -> + PidRefs = [spawn_monitor(fun() -> exit({pmap_res, Fun(Item)}) end) || Item <- Items], + lists:map( + fun({Pid, Ref}) -> + receive + {'DOWN', Ref, process, Pid, {pmap_res, Result}} -> Result; + {'DOWN', Ref, process, Pid, Error} -> exit(Error) + end + end, + PidRefs + ). + +chunk([], _N) -> + []; +chunk(Items, N) when length(Items) =< N -> + [Items]; +chunk(Items, N) -> + {Chunk, Rest} = lists:split(N, Items), + [Chunk | chunk(Rest, N)]. + group_info(DbName, DDocId, DbOptions) -> with_db(DbName, DbOptions, {couch_mrview, get_info, [DDocId]}). diff --git a/src/nouveau/src/nouveau_rpc.erl b/src/nouveau/src/nouveau_rpc.erl index 2037c7e7e..08d808a12 100644 --- a/src/nouveau/src/nouveau_rpc.erl +++ b/src/nouveau/src/nouveau_rpc.erl @@ -77,13 +77,5 @@ update_and_retry(DbName, Index, QueryArgs, UpdateLatency) -> rexi:reply(Else) end. -info(DbName, #index{} = Index0) -> - %% Incorporate the shard name into the record. - Index1 = Index0#index{dbname = DbName}, - case nouveau_api:index_info(Index1) of - {ok, Info0} -> - Info1 = Info0#{signature => Index0#index.sig}, - rexi:reply({ok, Info1}); - {error, Reason} -> - rexi:reply({error, Reason}) - end. +info(DbName, #index{} = Index) -> + rexi:reply(nouveau_util:index_info(DbName, Index)). diff --git a/src/nouveau/src/nouveau_util.erl b/src/nouveau/src/nouveau_util.erl index 567b77f49..e0b5700a4 100644 --- a/src/nouveau/src/nouveau_util.erl +++ b/src/nouveau/src/nouveau_util.erl @@ -22,6 +22,8 @@ index_name/1, design_doc_to_indexes/2, design_doc_to_index/3, + index_info/2, + index_info/3, verify_index_exists/2, ensure_local_purge_docs/2, maybe_create_local_purge_doc/2, @@ -40,6 +42,35 @@ index_name(#index{} = Index) -> node_prefix() -> atom_to_binary(node(), utf8). +index_info(DbName, #index{} = Index0) -> + %% Incorporate the shard name into the record. + Index1 = Index0#index{dbname = DbName}, + case nouveau_api:index_info(Index1) of + {ok, Info} -> + {ok, Info#{ + signature => Index1#index.sig, + pending_updates => pending_updates(DbName, Info) + }}; + {error, Reason} -> + {error, Reason} + end. + +index_info(DbName, DDoc, IndexName) -> + case design_doc_to_index(DbName, DDoc, IndexName) of + {ok, Index} -> + index_info(DbName, Index); + {error, Reason} -> + {error, Reason} + end. + +% Similar idea to view's pending_updates in couch_index. Return number of how +% much committed db seq is ahead of the index +pending_updates(DbName, #{} = Info) -> + GetCommSeq = fun(Db) -> couch_db:get_committed_update_seq(Db) end, + CommittedSeq = couch_util:with_db(DbName, GetCommSeq), + IndexSeq = maps:get(<<"update_seq">>, Info, 0), + max(CommittedSeq - IndexSeq, 0). + %% copied from dreyfus_index.erl design_doc_to_indexes(DbName, #doc{body = {Fields}} = Doc) -> RawIndexes = couch_util:get_value(<<"nouveau">>, Fields, {[]}),
