fix(integrations): address code-review findings

n8n: throw on FAILED memory events (was silently returned), poll on any
non-terminal status (PENDING or RUNNING), use top_k (not limit) for
search, coerce search/getAll responses to arrays, guard invalid-JSON
metadata, reject empty updates, mark user_id required; lint clean.

Zapier: coerce boolean fields (String(x) !== 'false') + boolean defaults
so "Wait"/"Infer = No" are honored, throw on HTTP >= 400 (was silent
empty/fake-success), throw on FAILED events, use top_k, numeric coercion
with clamping, mark user_id required.
This commit is contained in:
Himanshu-Sangshetti
2026-07-03 01:13:56 +08:00
parent 9882eea7b8
commit e66e97b725
8 changed files with 136 additions and 54 deletions
+6
View File
@@ -15,6 +15,12 @@ module.exports = {
files: ['./credentials/**/*.ts'],
plugins: ['eslint-plugin-n8n-nodes-base'],
extends: ['plugin:n8n-nodes-base/credentials'],
rules: {
// This rule only applies to nodes in n8n's main repository (where
// documentationUrl is an internal docs slug). Community nodes use a
// full external URL, so it is disabled here.
'n8n-nodes-base/cred-class-field-documentation-url-miscased': 'off',
},
},
{
files: ['./nodes/**/*.ts'],
@@ -8,7 +8,6 @@ import {
INodeTypeDescription,
JsonObject,
NodeApiError,
NodeConnectionTypes,
NodeOperationError,
sleep,
} from 'n8n-workflow';
@@ -31,8 +30,8 @@ export class Mem0 implements INodeType {
},
// Makes the node available to the AI Agent (Tools Agent) node.
usableAsTool: true,
inputs: [NodeConnectionTypes.Main],
outputs: [NodeConnectionTypes.Main],
inputs: ['main'],
outputs: ['main'],
credentials: [
{
name: 'mem0Api',
@@ -62,16 +61,10 @@ export class Mem0 implements INodeType {
description: 'Extract and store memories from messages',
},
{
name: 'Search',
value: 'search',
action: 'Search memories',
description: 'Semantic search over stored memories',
},
{
name: 'Get Many',
value: 'getAll',
action: 'Get many memories',
description: 'List stored memories for an entity',
name: 'Delete',
value: 'delete',
action: 'Delete a memory',
description: 'Delete a single memory by ID',
},
{
name: 'Get',
@@ -79,18 +72,24 @@ export class Mem0 implements INodeType {
action: 'Get a memory',
description: 'Retrieve a single memory by ID',
},
{
name: 'Get Many',
value: 'getAll',
action: 'Get many memories',
description: 'List stored memories for an entity',
},
{
name: 'Search',
value: 'search',
action: 'Search memories',
description: 'Semantic search over stored memories',
},
{
name: 'Update',
value: 'update',
action: 'Update a memory',
description: 'Update the text or metadata of a memory',
},
{
name: 'Delete',
value: 'delete',
action: 'Delete a memory',
description: 'Delete a single memory by ID',
},
],
default: 'add',
},
@@ -201,15 +200,16 @@ export class Mem0 implements INodeType {
name: 'userId',
type: 'string',
default: '',
required: true,
displayOptions: { show: { resource: ['memory'], operation: ['search'] } },
description: 'Restrict the search to this user',
description: 'Restrict the search to this user (required — the API needs an entity filter)',
},
{
displayName: 'Limit',
name: 'limit',
type: 'number',
typeOptions: { minValue: 1 },
default: 10,
default: 50,
displayOptions: { show: { resource: ['memory'], operation: ['search'] } },
description: 'Max number of results to return',
},
@@ -220,8 +220,9 @@ export class Mem0 implements INodeType {
name: 'userId',
type: 'string',
default: '',
required: true,
displayOptions: { show: { resource: ['memory'], operation: ['getAll'] } },
description: 'Restrict the listing to this user',
description: 'Restrict the listing to this user (required — the API needs an entity filter)',
},
{
displayName: 'Page',
@@ -316,31 +317,48 @@ export class Mem0 implements INodeType {
if (addFields.agent_id) body.agent_id = addFields.agent_id;
if (addFields.run_id) body.run_id = addFields.run_id;
if (addFields.metadata) {
body.metadata =
typeof addFields.metadata === 'string'
? JSON.parse(addFields.metadata as string)
: addFields.metadata;
try {
body.metadata =
typeof addFields.metadata === 'string'
? JSON.parse(addFields.metadata as string)
: addFields.metadata;
} catch {
throw new NodeOperationError(this.getNode(), 'Invalid JSON in "Metadata" field', {
itemIndex: i,
});
}
}
const addResp = await request('POST', '/v3/memories/add/', body);
const waitForCompletion = this.getNodeParameter('waitForCompletion', i, true) as boolean;
const addStatus = addResp.status as string | undefined;
const isTerminal = addStatus === 'SUCCEEDED' || addStatus === 'FAILED';
// infer=true returns {event_id, status:PENDING}; poll until done.
if (waitForCompletion && addResp.status === 'PENDING' && addResp.event_id) {
// infer=true returns {event_id, status:PENDING|RUNNING}; poll until terminal.
if (waitForCompletion && addResp.event_id && !isTerminal) {
responseData = await pollEvent(request, addResp.event_id as string, this);
} else if (addStatus === 'FAILED') {
throw new NodeOperationError(
this.getNode(),
`Mem0 memory add failed: ${(addResp.message as string) || 'unknown error'}`,
{ itemIndex: i },
);
} else {
responseData = addResp;
// Sync path (infer=false) returns {status:SUCCEEDED, results:[...]}; unwrap for consistency.
responseData = Array.isArray(addResp.results)
? (addResp.results as IDataObject[])
: addResp;
}
} else if (operation === 'search') {
const body: IDataObject = {
query: this.getNodeParameter('query', i) as string,
output_format: 'v1.1',
limit: this.getNodeParameter('limit', i, 10) as number,
top_k: this.getNodeParameter('limit', i, 50) as number,
};
const userId = this.getNodeParameter('userId', i, '') as string;
if (userId) body.filters = { user_id: userId };
const resp = await request('POST', '/v3/memories/search/', body);
responseData = (resp.results as IDataObject[]) ?? resp;
responseData = Array.isArray(resp.results) ? (resp.results as IDataObject[]) : [];
} else if (operation === 'getAll') {
const userId = this.getNodeParameter('userId', i, '') as string;
const page = this.getNodeParameter('page', i, 1) as number;
@@ -348,7 +366,7 @@ export class Mem0 implements INodeType {
const body: IDataObject = {};
if (userId) body.filters = { user_id: userId };
const resp = await request('POST', '/v3/memories/', body, { page, page_size: pageSize });
responseData = (resp.results as IDataObject[]) ?? resp;
responseData = Array.isArray(resp.results) ? (resp.results as IDataObject[]) : [];
} else if (operation === 'get') {
const memoryId = this.getNodeParameter('memoryId', i) as string;
responseData = await request('GET', `/v1/memories/${memoryId}/`);
@@ -358,7 +376,20 @@ export class Mem0 implements INodeType {
const text = this.getNodeParameter('text', i, '') as string;
const metadata = this.getNodeParameter('metadata', i, '') as string;
if (text) body.text = text;
if (metadata) body.metadata = typeof metadata === 'string' ? JSON.parse(metadata) : metadata;
if (metadata) {
try {
body.metadata = typeof metadata === 'string' ? JSON.parse(metadata) : metadata;
} catch {
throw new NodeOperationError(this.getNode(), 'Invalid JSON in "Metadata" field', {
itemIndex: i,
});
}
}
if (Object.keys(body).length === 0) {
throw new NodeOperationError(this.getNode(), 'Provide text or metadata to update', {
itemIndex: i,
});
}
responseData = await request('PUT', `/v1/memories/${memoryId}/`, body);
} else if (operation === 'delete') {
const memoryId = this.getNodeParameter('memoryId', i) as string;
@@ -392,9 +423,13 @@ async function pollEvent(
for (let attempt = 0; attempt < MAX_POLL_ATTEMPTS; attempt++) {
const event = await request('GET', `/v1/event/${eventId}/`);
const status = event.status as string;
if (status === 'SUCCEEDED' || status === 'FAILED') {
if (status === 'SUCCEEDED') {
return event;
}
if (status === 'FAILED') {
const reason = (event.error as string) || (event.message as string) || 'unknown error';
throw new NodeOperationError(ctx.getNode(), `Mem0 memory event ${eventId} failed: ${reason}`);
}
await sleep(POLL_INTERVAL_MS);
}
throw new NodeOperationError(
+1 -1
View File
@@ -27,5 +27,5 @@ module.exports = {
},
],
// Shown on the connection label in the Zap editor.
connectionLabel: 'Mem0 ({{bundle.authData.baseUrl}})',
connectionLabel: 'Mem0',
};
+32 -12
View File
@@ -8,27 +8,44 @@ const pollEvent = async (z, eventId) => {
for (let attempt = 0; attempt < MAX_POLL_ATTEMPTS; attempt++) {
const res = await z.request({ url: `/v1/event/${eventId}/`, method: 'GET' });
const status = res.data && res.data.status;
if (status === 'SUCCEEDED' || status === 'FAILED') {
if (status === 'SUCCEEDED') {
return res.data;
}
if (status === 'FAILED') {
const reason = (res.data && (res.data.error || res.data.message)) || 'unknown error';
throw new z.errors.Error(`Mem0 memory event ${eventId} failed: ${reason}`, 'Mem0EventFailed', 400);
}
await new Promise((resolve) => setTimeout(resolve, POLL_INTERVAL_MS));
}
throw new Error(`Timed out waiting for memory event ${eventId} to complete`);
throw new z.errors.Error(
`Timed out waiting for memory event ${eventId} to complete`,
'Mem0Timeout',
408,
);
};
const perform = async (z, bundle) => {
// Zapier boolean fields can arrive as the strings 'true'/'false'; coerce
// explicitly so "Infer = No" / "Wait = No" are honored.
const infer = String(bundle.inputData.infer) !== 'false';
const wait = String(bundle.inputData.waitForCompletion) !== 'false';
const body = {
messages: [{ role: bundle.inputData.role || 'user', content: bundle.inputData.content }],
infer: bundle.inputData.infer !== undefined ? bundle.inputData.infer : true,
infer,
};
if (bundle.inputData.user_id) body.user_id = bundle.inputData.user_id;
if (bundle.inputData.agent_id) body.agent_id = bundle.inputData.agent_id;
if (bundle.inputData.run_id) body.run_id = bundle.inputData.run_id;
if (bundle.inputData.metadata) {
body.metadata =
typeof bundle.inputData.metadata === 'string'
? JSON.parse(bundle.inputData.metadata)
: bundle.inputData.metadata;
try {
body.metadata =
typeof bundle.inputData.metadata === 'string'
? JSON.parse(bundle.inputData.metadata)
: bundle.inputData.metadata;
} catch (e) {
throw new z.errors.Error('Metadata must be valid JSON.', 'InvalidInput', 400);
}
}
const response = await z.request({
@@ -38,12 +55,15 @@ const perform = async (z, bundle) => {
});
const data = response.data;
const wait = bundle.inputData.waitForCompletion !== false;
// infer=true returns {event_id, status:PENDING}; optionally wait for the result.
if (wait && data.status && data.status !== 'SUCCEEDED' && data.event_id) {
// infer=true returns {event_id, status:PENDING|RUNNING}; optionally wait.
if (wait && data.event_id && data.status !== 'SUCCEEDED' && data.status !== 'FAILED') {
return pollEvent(z, data.event_id);
}
if (data.status === 'FAILED') {
const reason = data.error || data.message || 'unknown error';
throw new z.errors.Error(`Mem0 memory add failed: ${reason}`, 'Mem0EventFailed', 400);
}
return data;
};
@@ -78,14 +98,14 @@ module.exports = {
key: 'infer',
label: 'Infer',
type: 'boolean',
default: 'true',
default: true,
helpText: 'Run LLM extraction (async). Turn off to store verbatim (synchronous).',
},
{
key: 'waitForCompletion',
label: 'Wait for Completion',
type: 'boolean',
default: 'true',
default: true,
helpText: 'Poll until extraction finishes and return the resulting memories.',
},
],
+12 -1
View File
@@ -14,7 +14,9 @@ const includeApiKey = (request, z, bundle) => {
return request;
};
// Surface auth failures as Zapier auth errors so users are prompted to reconnect.
// Surface HTTP failures as errors. z.request does NOT throw on non-2xx by
// default, so without this a 4xx/5xx would flow downstream as a fake success
// (empty search results / error body returned as a created memory).
const handleBadResponses = (response, z, _bundle) => {
if (response.status === 401 || response.status === 403) {
throw new z.errors.Error(
@@ -23,6 +25,15 @@ const handleBadResponses = (response, z, _bundle) => {
response.status,
);
}
if (response.status >= 400) {
const data = response.data || {};
const detail = data.detail || data.error || data.message || response.content || 'unknown error';
throw new z.errors.Error(
`Mem0 API request failed (HTTP ${response.status}): ${detail}`,
'Mem0ApiError',
response.status,
);
}
return response;
};
@@ -7,7 +7,7 @@ const perform = async (z, bundle) => {
const response = await z.request({
url: '/v3/memories/',
method: 'POST',
params: { page: 1, page_size: bundle.inputData.limit || 50 },
params: { page: 1, page_size: Math.max(1, Math.floor(Number(bundle.inputData.limit) || 50)) },
body,
});
const data = response.data;
@@ -24,8 +24,8 @@ module.exports = {
operation: {
perform,
inputFields: [
{ key: 'user_id', label: 'User ID', type: 'string' },
{ key: 'limit', label: 'Limit', type: 'integer', default: '50' },
{ key: 'user_id', label: 'User ID', type: 'string', required: true },
{ key: 'limit', label: 'Limit', type: 'integer', default: 50 },
],
sample: { id: '00000000-0000-0000-0000-000000000000', memory: 'User loves hiking' },
},
@@ -4,7 +4,7 @@ const perform = async (z, bundle) => {
const body = {
query: bundle.inputData.query,
output_format: 'v1.1',
limit: bundle.inputData.limit || 10,
top_k: Math.max(1, Math.floor(Number(bundle.inputData.limit) || 10)),
};
if (bundle.inputData.user_id) body.filters = { user_id: bundle.inputData.user_id };
@@ -29,8 +29,8 @@ module.exports = {
perform,
inputFields: [
{ key: 'query', label: 'Query', type: 'string', required: true },
{ key: 'user_id', label: 'User ID', type: 'string' },
{ key: 'limit', label: 'Limit', type: 'integer', default: '10' },
{ key: 'user_id', label: 'User ID', type: 'string', required: true },
{ key: 'limit', label: 'Limit', type: 'integer', default: 10 },
],
sample: { id: '00000000-0000-0000-0000-000000000000', memory: 'User loves hiking' },
},
@@ -66,4 +66,14 @@ describe('Mem0 Zapier integration (E2E)', () => {
});
expect(afterDelete.length).toBe(0);
});
it('surfaces API errors instead of returning an empty array (search needs a filter)', async () => {
// filters is required by the API; omitting it must throw, not return [].
await expect(
appTester(App.searches.search_memories.operation.perform, {
authData,
inputData: { query: 'anything' },
}),
).rejects.toThrow();
});
});