fix(gateway,api): reach NATS over IPv4 and survive boot races (#3189)

This commit is contained in:
Hampus
2026-10-04 03:06:12 +02:00
committed by GitHub
parent 1544e58e76
commit f3c777b244
6 changed files with 75 additions and 5 deletions
@@ -17,6 +17,16 @@ export interface MeilisearchTask {
}; };
} }
export class MeilisearchTaskError extends Error {
readonly code: string | undefined;
constructor(message: string, code: string | undefined) {
super(message);
this.name = 'MeilisearchTaskError';
this.code = code;
}
}
export interface MeilisearchClient { export interface MeilisearchClient {
request<TResponse>(method: string, path: string, body?: unknown): Promise<TResponse>; request<TResponse>(method: string, path: string, body?: unknown): Promise<TResponse>;
waitForTask(taskUid: number): Promise<void>; waitForTask(taskUid: number): Promise<void>;
@@ -84,7 +94,10 @@ export class MeilisearchHttpClient implements MeilisearchClient {
return; return;
} }
if (task.status === 'failed' || task.status === 'canceled') { if (task.status === 'failed' || task.status === 'canceled') {
throw new Error(task.error?.message ?? `Meilisearch task ${taskUid} ${task.status}`); throw new MeilisearchTaskError(
task.error?.message ?? `Meilisearch task ${taskUid} ${task.status}`,
task.error?.code,
);
} }
await new Promise((resolve) => setTimeout(resolve, TASK_POLL_INTERVAL_MS)); await new Promise((resolve) => setTimeout(resolve, TASK_POLL_INTERVAL_MS));
} }
@@ -1,6 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later // SPDX-License-Identifier: AGPL-3.0-or-later
import type {MeilisearchClient, MeilisearchTask} from '@app/api/search/meilisearch/MeilisearchClient'; import type {MeilisearchClient, MeilisearchTask} from '@app/api/search/meilisearch/MeilisearchClient';
import {MeilisearchTaskError} from '@app/api/search/meilisearch/MeilisearchClient';
import {MeilisearchMessageAdapter} from '@app/api/search/meilisearch/MeilisearchDomainAdapters'; import {MeilisearchMessageAdapter} from '@app/api/search/meilisearch/MeilisearchDomainAdapters';
import {MEILISEARCH_MAX_TRACKED_BULK_TASKS} from '@app/api/search/meilisearch/MeilisearchIndexAdapter'; import {MEILISEARCH_MAX_TRACKED_BULK_TASKS} from '@app/api/search/meilisearch/MeilisearchIndexAdapter';
import type {SearchableMessage} from '@fluxer/schema/src/contracts/search/SearchDocumentTypes'; import type {SearchableMessage} from '@fluxer/schema/src/contracts/search/SearchDocumentTypes';
@@ -15,6 +16,7 @@ interface RecordedMeilisearchRequest {
class FakeMeilisearchClient implements MeilisearchClient { class FakeMeilisearchClient implements MeilisearchClient {
readonly requests: Array<RecordedMeilisearchRequest> = []; readonly requests: Array<RecordedMeilisearchRequest> = [];
readonly waitedTaskUids: Array<number> = []; readonly waitedTaskUids: Array<number> = [];
readonly failedTasks = new Map<number, MeilisearchTaskError>();
private nextTaskUid = 1; private nextTaskUid = 1;
indexExists = false; indexExists = false;
@@ -50,6 +52,10 @@ class FakeMeilisearchClient implements MeilisearchClient {
async waitForTask(taskUid: number): Promise<void> { async waitForTask(taskUid: number): Promise<void> {
this.waitedTaskUids.push(taskUid); this.waitedTaskUids.push(taskUid);
const failure = this.failedTasks.get(taskUid);
if (failure) {
throw failure;
}
} }
clear(): void { clear(): void {
@@ -87,6 +93,26 @@ describe('MeilisearchMessageAdapter', () => {
}); });
}); });
it('treats an index created concurrently by another process as created', async () => {
const client = new FakeMeilisearchClient();
client.failedTasks.set(1, new MeilisearchTaskError('Index `messages` already exists.', 'index_already_exists'));
const adapter = new MeilisearchMessageAdapter({client});
await adapter.initialize();
expect(adapter.isAvailable()).toBe(true);
expect(client.waitedTaskUids).toEqual([1, 2, 3, 4, 5]);
});
it('still fails when creating the index fails for another reason', async () => {
const client = new FakeMeilisearchClient();
client.failedTasks.set(1, new MeilisearchTaskError('Index uid is invalid.', 'invalid_index_uid'));
const adapter = new MeilisearchMessageAdapter({client});
await expect(adapter.initialize()).rejects.toThrow('Index uid is invalid.');
expect(adapter.isAvailable()).toBe(false);
});
it('builds Meilisearch search requests from message filters', async () => { it('builds Meilisearch search requests from message filters', async () => {
const client = new FakeMeilisearchClient(); const client = new FakeMeilisearchClient();
client.indexExists = true; client.indexExists = true;
@@ -1,6 +1,7 @@
// SPDX-License-Identifier: AGPL-3.0-or-later // SPDX-License-Identifier: AGPL-3.0-or-later
import type {MeilisearchClient, MeilisearchTask} from '@app/api/search/meilisearch/MeilisearchClient'; import type {MeilisearchClient, MeilisearchTask} from '@app/api/search/meilisearch/MeilisearchClient';
import {MeilisearchTaskError} from '@app/api/search/meilisearch/MeilisearchClient';
import type {MeilisearchFilter} from '@app/api/search/meilisearch/MeilisearchFilterUtils'; import type {MeilisearchFilter} from '@app/api/search/meilisearch/MeilisearchFilterUtils';
import {joinMeiliFilters} from '@app/api/search/meilisearch/MeilisearchFilterUtils'; import {joinMeiliFilters} from '@app/api/search/meilisearch/MeilisearchFilterUtils';
import type {MeilisearchIndexDefinition} from '@app/api/search/meilisearch/MeilisearchIndexDefinitions'; import type {MeilisearchIndexDefinition} from '@app/api/search/meilisearch/MeilisearchIndexDefinitions';
@@ -58,7 +59,13 @@ export class MeilisearchIndexAdapter<
uid, uid,
primaryKey: this.indexDefinition.primaryKey, primaryKey: this.indexDefinition.primaryKey,
}); });
await this.client.waitForTask(task.taskUid); try {
await this.client.waitForTask(task.taskUid);
} catch (error) {
if (!(error instanceof MeilisearchTaskError && error.code === 'index_already_exists')) {
throw error;
}
}
} }
await Promise.all([ await Promise.all([
this.applySetting('PUT', 'searchable-attributes', this.indexDefinition.searchableAttributes), this.applySetting('PUT', 'searchable-attributes', this.indexDefinition.searchableAttributes),
@@ -414,7 +414,8 @@ export class JetStreamWorkerQueue {
this.requireConsumerConfiguration(existing, lane.consumerName); this.requireConsumerConfiguration(existing, lane.consumerName);
const updated = await jsm.consumers.update(STREAM_NAME, lane.consumerName, config); const updated = await jsm.consumers.update(STREAM_NAME, lane.consumerName, config);
this.requireConsumerConfiguration(updated, lane.consumerName); this.requireConsumerConfiguration(updated, lane.consumerName);
if (updated.created !== existing.created) { const current = await this.readConsumer(jsm, lane.consumerName);
if (current?.created !== existing.created) {
throw new Error(`Worker consumer ${lane.consumerName} was replaced during startup`); throw new Error(`Worker consumer ${lane.consumerName} was replaced during startup`);
} }
Logger.info({lane: lane.name, consumer: lane.consumerName}, 'Consumer updated without resetting delivery state'); Logger.info({lane: lane.name, consumer: lane.consumerName}, 'Consumer updated without resetting delivery state');
@@ -5,6 +5,7 @@
-export([ -export([
get_pool_conn/0, get_pool_conn/0,
connect/3,
connect_slot/2, connect_slot/2,
handle_conn_down/4, handle_conn_down/4,
find_slot_by_conn/2, find_slot_by_conn/2,
@@ -100,7 +101,7 @@ spawn_connect_slot(Idx, Host, Port, Opts, Parent) ->
ok. ok.
connect_slot_and_notify(Idx, Host, Port, Opts, Parent) -> connect_slot_and_notify(Idx, Host, Port, Opts, Parent) ->
_ = _ =
case nats:connect(Host, Port, Opts) of case connect(Host, Port, Opts) of
{ok, Conn} -> {ok, Conn} ->
ok = nats:controlling_process(Conn, Parent), ok = nats:controlling_process(Conn, Parent),
Parent ! {nats_pool_connect_result, Idx, self(), {ok, Conn}}; Parent ! {nats_pool_connect_result, Idx, self(), {ok, Conn}};
@@ -109,6 +110,21 @@ connect_slot_and_notify(Idx, Host, Port, Opts, Parent) ->
end, end,
ok. ok.
-spec connect(string(), inet:port_number(), map()) ->
{ok, nats:conn()} | ignore | {error, term()}.
connect(Host, Port, Opts) ->
Server = #{
scheme => <<"nats">>, host => Host, port => Port, family => address_family(Host)
},
nats:connect(Server, Opts).
-spec address_family(string()) -> inet | inet6.
address_family(Host) ->
case inet:getaddr(Host, inet) of
{ok, _} -> inet;
{error, _} -> inet6
end.
-spec handle_conn_down(non_neg_integer(), nats:conn(), pos_integer(), map()) -> map(). -spec handle_conn_down(non_neg_integer(), nats:conn(), pos_integer(), map()) -> map().
handle_conn_down( handle_conn_down(
Idx, Idx,
@@ -385,6 +401,13 @@ build_connect_opts_with_token_test() ->
#{auth_token => <<"secret">>, buffer_size => 0}, build_connect_opts(<<"secret">>) #{auth_token => <<"secret">>, buffer_size => 0}, build_connect_opts(<<"secret">>)
). ).
address_family_prefers_ipv4_test() ->
?assertEqual(inet, address_family("127.0.0.1")),
?assertEqual(inet, address_family("localhost")).
address_family_falls_back_to_ipv6_test() ->
?assertEqual(inet6, address_family("::1")).
build_connect_opts_empty_test() -> build_connect_opts_empty_test() ->
?assertEqual(#{buffer_size => 0}, build_connect_opts(undefined)). ?assertEqual(#{buffer_size => 0}, build_connect_opts(undefined)).
@@ -152,7 +152,7 @@ spawn_connect(Host, Port, AuthToken, State) ->
-spec connect_and_notify(string(), inet:port_number(), map(), pid()) -> ok. -spec connect_and_notify(string(), inet:port_number(), map(), pid()) -> ok.
connect_and_notify(Host, Port, Opts, Parent) -> connect_and_notify(Host, Port, Opts, Parent) ->
_ = _ =
case nats:connect(Host, Port, Opts) of case gateway_nats_pool_conn:connect(Host, Port, Opts) of
{ok, Conn} -> {ok, Conn} ->
ok = nats:controlling_process(Conn, Parent), ok = nats:controlling_process(Conn, Parent),
Parent ! {nats_connect_result, self(), {ok, Conn}}; Parent ! {nats_connect_result, self(), {ok, Conn}};