feat(core): Add poller_state table for durable poll trigger cursors (no-changelog) (#35128)

Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
Co-authored-by: Lorent Lempereur <lorent.lempereur@n8n.io>
This commit is contained in:
Emilia
2026-08-11 17:31:00 +01:00
committed by GitHub
parent a447ac422e
commit 93ff542e92
12 changed files with 401 additions and 2 deletions
+11
View File
@@ -92,6 +92,7 @@ Auto-generated from the PostgreSQL migrations in @n8n/db. Do not edit by hand.
| [public.oauth_clients](public.oauth_clients.md) | 10 | | BASE TABLE |
| [public.oauth_refresh_tokens](public.oauth_refresh_tokens.md) | 7 | | BASE TABLE |
| [public.oauth_user_consents](public.oauth_user_consents.md) | 5 | | BASE TABLE |
| [public.poller_state](public.poller_state.md) | 7 | | BASE TABLE |
| [public.processed_data](public.processed_data.md) | 5 | | BASE TABLE |
| [public.project](public.project.md) | 9 | | BASE TABLE |
| [public.project_relation](public.project_relation.md) | 5 | | BASE TABLE |
@@ -280,6 +281,7 @@ erDiagram
"public.oauth_refresh_tokens" }o--|| "public.oauth_clients" : "FOREIGN KEY (#quot;clientId#quot;) REFERENCES oauth_clients(id) ON DELETE CASCADE"
"public.oauth_user_consents" }o--|| "public.user" : "FOREIGN KEY (#quot;userId#quot;) REFERENCES #quot;user#quot;(id) ON DELETE CASCADE"
"public.oauth_user_consents" }o--|| "public.oauth_clients" : "FOREIGN KEY (#quot;clientId#quot;) REFERENCES oauth_clients(id) ON DELETE CASCADE"
"public.poller_state" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.processed_data" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.project" }o--o| "public.user" : "FOREIGN KEY (#quot;creatorId#quot;) REFERENCES #quot;user#quot;(id) ON DELETE SET NULL"
"public.project_relation" }o--|| "public.user" : "FOREIGN KEY (#quot;userId#quot;) REFERENCES #quot;user#quot;(id) ON DELETE CASCADE"
@@ -1180,6 +1182,15 @@ erDiagram
json scope
uuid userId FK
}
"public.poller_state" {
timestamp_3__with_time_zone backoffUntil
integer consecutiveErrors
timestamp_3__with_time_zone createdAt
json cursor
varchar_36_ nodeId
timestamp_3__with_time_zone updatedAt
varchar_36_ workflowId FK
}
"public.processed_data" {
varchar_255_ context
timestamp_3__with_time_zone createdAt
@@ -0,0 +1,76 @@
# public.poller_state
## Columns
| Name | Type | Default | Nullable | Children | Parents | Comment |
| ---- | ---- | ------- | -------- | -------- | ------- | ------- |
| backoffUntil | timestamp(3) with time zone | | true | | | Time before which no poll is attempted; NULL when not backing off. |
| consecutiveErrors | integer | 0 | false | | | Polls that have failed since the last successful one. |
| createdAt | timestamp(3) with time zone | CURRENT_TIMESTAMP(3) | false | | | |
| cursor | json | '{}'::json | false | | | How far the poll node has consumed its source, in whatever shape that node uses. |
| nodeId | varchar(36) | | false | | | |
| updatedAt | timestamp(3) with time zone | CURRENT_TIMESTAMP(3) | false | | | |
| workflowId | varchar(36) | | false | | [public.workflow_entity](public.workflow_entity.md) | |
## Constraints
| Name | Type | Definition |
| ---- | ---- | ---------- |
| FK_poller_state_workflowId | FOREIGN KEY | FOREIGN KEY ("workflowId") REFERENCES workflow_entity(id) ON DELETE CASCADE |
| PK_eda59336a5dd1e6a12f4d4b7287 | PRIMARY KEY | PRIMARY KEY ("workflowId", "nodeId") |
| poller_state_consecutiveErrors_not_null | n | NOT NULL "consecutiveErrors" |
| poller_state_createdAt_not_null | n | NOT NULL "createdAt" |
| poller_state_cursor_not_null | n | NOT NULL cursor |
| poller_state_nodeId_not_null | n | NOT NULL "nodeId" |
| poller_state_updatedAt_not_null | n | NOT NULL "updatedAt" |
| poller_state_workflowId_not_null | n | NOT NULL "workflowId" |
## Indexes
| Name | Definition |
| ---- | ---------- |
| PK_eda59336a5dd1e6a12f4d4b7287 | CREATE UNIQUE INDEX "PK_eda59336a5dd1e6a12f4d4b7287" ON public.poller_state USING btree ("workflowId", "nodeId") |
## Relations
```mermaid
erDiagram
"public.poller_state" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.poller_state" {
timestamp_3__with_time_zone backoffUntil
integer consecutiveErrors
timestamp_3__with_time_zone createdAt
json cursor
varchar_36_ nodeId
timestamp_3__with_time_zone updatedAt
varchar_36_ workflowId FK
}
"public.workflow_entity" {
boolean active
varchar_36_ activeVersionId FK
json connections
timestamp_3__with_time_zone createdAt
text description
varchar_36_ id
boolean isArchived
json meta
varchar_128_ name
json nodeGroups
json nodes
varchar_36_ parentFolderId FK
json pinData
json settings
varchar sourceWorkflowId
json staticData
integer triggerCount
timestamp_3__with_time_zone updatedAt
integer versionCounter
character_36_ versionId
}
```
---
> Generated by [tbls](https://github.com/k1LoW/tbls)
@@ -9,7 +9,7 @@
| connections | json | | false | | | |
| createdAt | timestamp(3) with time zone | CURRENT_TIMESTAMP(3) | false | | | |
| description | text | | true | | | |
| id | varchar(36) | | false | [public.ai_builder_temporary_workflow](public.ai_builder_temporary_workflow.md) [public.chat_hub_messages](public.chat_hub_messages.md) [public.chat_hub_sessions](public.chat_hub_sessions.md) [public.evaluation_collection](public.evaluation_collection.md) [public.evaluation_config](public.evaluation_config.md) [public.execution_entity](public.execution_entity.md) [public.insights_metadata](public.insights_metadata.md) [public.processed_data](public.processed_data.md) [public.shared_workflow](public.shared_workflow.md) [public.test_run](public.test_run.md) [public.webhook_entity](public.webhook_entity.md) [public.workflow_builder_session](public.workflow_builder_session.md) [public.workflow_dependency](public.workflow_dependency.md) [public.workflow_history](public.workflow_history.md) [public.workflow_publication_trigger_status](public.workflow_publication_trigger_status.md) [public.workflow_publish_history](public.workflow_publish_history.md) [public.workflow_published_version](public.workflow_published_version.md) [public.workflow_review_activity](public.workflow_review_activity.md) [public.workflow_review_request_workflow](public.workflow_review_request_workflow.md) [public.workflows_tags](public.workflows_tags.md) | | |
| id | varchar(36) | | false | [public.ai_builder_temporary_workflow](public.ai_builder_temporary_workflow.md) [public.chat_hub_messages](public.chat_hub_messages.md) [public.chat_hub_sessions](public.chat_hub_sessions.md) [public.evaluation_collection](public.evaluation_collection.md) [public.evaluation_config](public.evaluation_config.md) [public.execution_entity](public.execution_entity.md) [public.insights_metadata](public.insights_metadata.md) [public.poller_state](public.poller_state.md) [public.processed_data](public.processed_data.md) [public.shared_workflow](public.shared_workflow.md) [public.test_run](public.test_run.md) [public.webhook_entity](public.webhook_entity.md) [public.workflow_builder_session](public.workflow_builder_session.md) [public.workflow_dependency](public.workflow_dependency.md) [public.workflow_history](public.workflow_history.md) [public.workflow_publication_trigger_status](public.workflow_publication_trigger_status.md) [public.workflow_publish_history](public.workflow_publish_history.md) [public.workflow_published_version](public.workflow_published_version.md) [public.workflow_review_activity](public.workflow_review_activity.md) [public.workflow_review_request_workflow](public.workflow_review_request_workflow.md) [public.workflows_tags](public.workflows_tags.md) | | |
| isArchived | boolean | false | false | | | |
| meta | json | | true | | | |
| name | varchar(128) | | false | | | |
@@ -73,6 +73,7 @@ erDiagram
"public.evaluation_config" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.execution_entity" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.insights_metadata" }o--o| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE SET NULL"
"public.poller_state" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.processed_data" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.shared_workflow" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
"public.test_run" }o--|| "public.workflow_entity" : "FOREIGN KEY (#quot;workflowId#quot;) REFERENCES workflow_entity(id) ON DELETE CASCADE"
@@ -216,6 +217,15 @@ erDiagram
varchar_36_ workflowId FK
varchar_128_ workflowName
}
"public.poller_state" {
timestamp_3__with_time_zone backoffUntil
integer consecutiveErrors
timestamp_3__with_time_zone createdAt
json cursor
varchar_36_ nodeId
timestamp_3__with_time_zone updatedAt
varchar_36_ workflowId FK
}
"public.processed_data" {
varchar_255_ context
timestamp_3__with_time_zone createdAt
+11
View File
@@ -92,6 +92,7 @@ Auto-generated from the SQLite migrations in @n8n/db. Do not edit by hand.
| [oauth_clients](oauth_clients.md) | 10 | | table |
| [oauth_refresh_tokens](oauth_refresh_tokens.md) | 7 | | table |
| [oauth_user_consents](oauth_user_consents.md) | 5 | | table |
| [poller_state](poller_state.md) | 7 | | table |
| [processed_data](processed_data.md) | 5 | | table |
| [project](project.md) | 9 | | table |
| [project_relation](project_relation.md) | 5 | | table |
@@ -263,6 +264,7 @@ erDiagram
"oauth_refresh_tokens" }o--|| "user" : "FOREIGN KEY (userId) REFERENCES user (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"oauth_user_consents" }o--|| "user" : "FOREIGN KEY (userId) REFERENCES user (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"oauth_user_consents" }o--|| "oauth_clients" : "FOREIGN KEY (clientId) REFERENCES oauth_clients (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"poller_state" |o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"processed_data" |o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"project" }o--o| "user" : "FOREIGN KEY (creatorId) REFERENCES user (id) ON UPDATE NO ACTION ON DELETE SET NULL MATCH NONE"
"project_relation" }o--|| "role" : "FOREIGN KEY (role) REFERENCES role (slug) ON UPDATE NO ACTION ON DELETE NO ACTION MATCH NONE"
@@ -1167,6 +1169,15 @@ erDiagram
TEXT scope
varchar userId FK
}
"poller_state" {
datetime_3_ backoffUntil
INTEGER consecutiveErrors
datetime_3_ createdAt
TEXT cursor
varchar_36_ nodeId PK
datetime_3_ updatedAt
varchar_36_ workflowId PK
}
"processed_data" {
varchar_255_ context PK
datetime_3_ createdAt
@@ -0,0 +1,83 @@
# poller_state
## Description
<details>
<summary><strong>Table Definition</strong></summary>
```sql
CREATE TABLE "poller_state" ("workflowId" varchar(36) NOT NULL, "nodeId" varchar(36) NOT NULL, "cursor" text NOT NULL DEFAULT ('{}'), "consecutiveErrors" integer NOT NULL DEFAULT (0), "backoffUntil" datetime(3), "createdAt" datetime(3) NOT NULL DEFAULT (STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW')), "updatedAt" datetime(3) NOT NULL DEFAULT (STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW')), CONSTRAINT "FK_poller_state_workflowId" FOREIGN KEY ("workflowId") REFERENCES "workflow_entity" ("id") ON DELETE CASCADE, PRIMARY KEY ("workflowId", "nodeId"))
```
</details>
## Columns
| Name | Type | Default | Nullable | Children | Parents | Comment |
| ---- | ---- | ------- | -------- | -------- | ------- | ------- |
| backoffUntil | datetime(3) | | true | | | |
| consecutiveErrors | INTEGER | 0 | false | | | |
| createdAt | datetime(3) | STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW') | false | | | |
| cursor | TEXT | '{}' | false | | | |
| nodeId | varchar(36) | | false | | | |
| updatedAt | datetime(3) | STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW') | false | | | |
| workflowId | varchar(36) | | false | | [workflow_entity](workflow_entity.md) | |
## Constraints
| Name | Type | Definition |
| ---- | ---- | ---------- |
| - (Foreign key ID: 0) | FOREIGN KEY | FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE |
| nodeId | PRIMARY KEY | PRIMARY KEY (nodeId) |
| sqlite_autoindex_poller_state_1 | PRIMARY KEY | PRIMARY KEY (workflowId, nodeId) |
| workflowId | PRIMARY KEY | PRIMARY KEY (workflowId) |
## Indexes
| Name | Definition |
| ---- | ---------- |
| sqlite_autoindex_poller_state_1 | PRIMARY KEY (workflowId, nodeId) |
## Relations
```mermaid
erDiagram
"poller_state" |o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"poller_state" {
datetime_3_ backoffUntil
INTEGER consecutiveErrors
datetime_3_ createdAt
TEXT cursor
varchar_36_ nodeId PK
datetime_3_ updatedAt
varchar_36_ workflowId PK
}
"workflow_entity" {
boolean active
varchar_36_ activeVersionId FK
TEXT connections
datetime_3_ createdAt
TEXT description
varchar_36_ id PK
boolean isArchived
TEXT meta
varchar_128_ name
TEXT nodeGroups
TEXT nodes
varchar_36_ parentFolderId FK
TEXT pinData
TEXT settings
varchar sourceWorkflowId
TEXT staticData
INTEGER triggerCount
datetime_3_ updatedAt
INTEGER versionCounter
varchar_36_ versionId
}
```
---
> Generated by [tbls](https://github.com/k1LoW/tbls)
@@ -20,7 +20,7 @@ CREATE TABLE "workflow_entity" ("id" varchar(36) PRIMARY KEY NOT NULL, "name" va
| connections | TEXT | | true | | | |
| createdAt | datetime(3) | STRFTIME('%Y-%m-%d %H:%M:%f', 'NOW') | false | | | |
| description | TEXT | | true | | | |
| id | varchar(36) | | false | [ai_builder_temporary_workflow](ai_builder_temporary_workflow.md) [chat_hub_messages](chat_hub_messages.md) [chat_hub_sessions](chat_hub_sessions.md) [evaluation_collection](evaluation_collection.md) [evaluation_config](evaluation_config.md) [execution_entity](execution_entity.md) [insights_metadata](insights_metadata.md) [processed_data](processed_data.md) [shared_workflow](shared_workflow.md) [test_run](test_run.md) [workflow_builder_session](workflow_builder_session.md) [workflow_dependency](workflow_dependency.md) [workflow_history](workflow_history.md) [workflow_publication_trigger_status](workflow_publication_trigger_status.md) [workflow_publish_history](workflow_publish_history.md) [workflow_published_version](workflow_published_version.md) [workflow_review_activity](workflow_review_activity.md) [workflow_review_request_workflow](workflow_review_request_workflow.md) [workflows_tags](workflows_tags.md) | | |
| id | varchar(36) | | false | [ai_builder_temporary_workflow](ai_builder_temporary_workflow.md) [chat_hub_messages](chat_hub_messages.md) [chat_hub_sessions](chat_hub_sessions.md) [evaluation_collection](evaluation_collection.md) [evaluation_config](evaluation_config.md) [execution_entity](execution_entity.md) [insights_metadata](insights_metadata.md) [poller_state](poller_state.md) [processed_data](processed_data.md) [shared_workflow](shared_workflow.md) [test_run](test_run.md) [workflow_builder_session](workflow_builder_session.md) [workflow_dependency](workflow_dependency.md) [workflow_history](workflow_history.md) [workflow_publication_trigger_status](workflow_publication_trigger_status.md) [workflow_publish_history](workflow_publish_history.md) [workflow_published_version](workflow_published_version.md) [workflow_review_activity](workflow_review_activity.md) [workflow_review_request_workflow](workflow_review_request_workflow.md) [workflows_tags](workflows_tags.md) | | |
| isArchived | boolean | FALSE | false | | | |
| meta | TEXT | | true | | | |
| name | varchar(128) | | false | | | |
@@ -72,6 +72,7 @@ erDiagram
"evaluation_config" }o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"execution_entity" }o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"insights_metadata" }o--o| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE SET NULL MATCH NONE"
"poller_state" |o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"processed_data" |o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"shared_workflow" |o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
"test_run" }o--|| "workflow_entity" : "FOREIGN KEY (workflowId) REFERENCES workflow_entity (id) ON UPDATE NO ACTION ON DELETE CASCADE MATCH NONE"
@@ -215,6 +216,15 @@ erDiagram
varchar_16_ workflowId FK
varchar_128_ workflowName
}
"poller_state" {
datetime_3_ backoffUntil
INTEGER consecutiveErrors
datetime_3_ createdAt
TEXT cursor
varchar_36_ nodeId PK
datetime_3_ updatedAt
varchar_36_ workflowId PK
}
"processed_data" {
varchar_255_ context PK
datetime_3_ createdAt
+3
View File
@@ -26,6 +26,7 @@ import { Folder } from './folder';
import { FolderTagMapping } from './folder-tag-mapping';
import { InstanceCredentialAssignment } from './instance-credential-assignment';
import { InvalidAuthToken } from './invalid-auth-token';
import { PollerState } from './poller-state';
import { ProcessedData } from './processed-data';
import { Project } from './project';
import { ProjectRelation } from './project-relation';
@@ -99,6 +100,7 @@ export {
InvalidAuthToken,
InstanceCredentialAssignment,
AiBuilderTemporaryWorkflow,
PollerState,
ProcessedData,
Settings,
Variables,
@@ -181,6 +183,7 @@ export const entities = {
InvalidAuthToken,
InstanceCredentialAssignment,
AiBuilderTemporaryWorkflow,
PollerState,
ProcessedData,
Settings,
Variables,
@@ -0,0 +1,34 @@
import { Column, Entity, PrimaryColumn } from '@n8n/typeorm';
import { DateTimeColumn, JsonColumn, WithTimestamps } from './abstract-entity';
export type PollerCursor = Record<string, unknown>;
/**
* Durable state for one poll trigger node.
*
* State is kept per node rather than per workflow so that two poll nodes in one
* workflow don't contend on every write.
*/
@Entity({ name: 'poller_state' })
export class PollerState extends WithTimestamps {
@PrimaryColumn({ type: 'varchar', length: 36 })
workflowId: string;
@PrimaryColumn({ type: 'varchar', length: 36 })
nodeId: string;
/**
* How far the node has consumed its source. The shape is the node's own,
* e.g. a timestamp, a page token, or a list of already-emitted ids.
*/
@JsonColumn({ default: '{}' })
cursor: PollerCursor;
@Column({ type: 'int', default: 0 })
consecutiveErrors: number;
/** Time before which no poll is attempted. */
@DateTimeColumn({ nullable: true })
backoffUntil: Date | null;
}
@@ -0,0 +1,32 @@
import type { MigrationContext, ReversibleMigration } from '../migration-types';
export class CreatePollerStateTable1785926660580 implements ReversibleMigration {
async up({ schemaBuilder: { createTable, column }, tablePrefix }: MigrationContext) {
await createTable('poller_state')
.withColumns(
column('workflowId').varchar(36).primary,
column('nodeId').varchar(36).primary,
column('cursor')
.json.notNull.default("'{}'")
.comment(
'How far the poll node has consumed its source, in whatever shape that node uses.',
),
column('consecutiveErrors')
.int.notNull.default(0)
.comment('Polls that have failed since the last successful one.'),
column('backoffUntil')
.timestampTimezone()
.comment('Time before which no poll is attempted; NULL when not backing off.'),
)
.withTimestamps.withForeignKey('workflowId', {
tableName: 'workflow_entity',
columnName: 'id',
onDelete: 'CASCADE',
name: `FK_${tablePrefix}poller_state_workflowId`,
});
}
async down({ schemaBuilder: { dropTable } }: MigrationContext) {
await dropTable('poller_state');
}
}
@@ -54,6 +54,7 @@ export type {
DeleteFinishedTasksOptions,
ScheduledTaskMetricSnapshot,
} from './scheduled-task.repository';
export { PollerStateRepository, type PollerCursor } from './poller-state.repository';
export { ProcessedDataRepository } from './processed-data.repository';
export { SettingsRepository } from './settings.repository';
export { TagRepository } from './tag.repository';
@@ -0,0 +1,29 @@
import { Service } from '@n8n/di';
import { DataSource } from '@n8n/typeorm';
import { PollerState } from '../entities';
import { BaseRepository } from './base-repository';
import type { PollerCursor } from '../entities/poller-state';
import type { OperationContext } from '../services/transaction';
export type { PollerCursor } from '../entities/poller-state';
@Service()
export class PollerStateRepository extends BaseRepository<PollerState> {
constructor(dataSource: DataSource) {
super(PollerState, dataSource.manager);
}
/** The node's stored cursor, or `null` if it has never polled. */
async findCursor(
workflowId: string,
nodeId: string,
ctx: OperationContext = {},
): Promise<PollerCursor | null> {
const row = await this.managerFor(ctx).findOne(PollerState, {
select: ['cursor'],
where: { workflowId, nodeId },
});
return row === null ? null : row.cursor;
}
}
@@ -0,0 +1,99 @@
import { createWorkflow, testDb } from '@n8n/backend-test-utils';
import { type PollerCursor, PollerStateRepository, WorkflowRepository } from '@n8n/db';
import { Container } from '@n8n/di';
describe('PollerStateRepository', () => {
let repository: PollerStateRepository;
let workflowRepository: WorkflowRepository;
let workflowId: string;
beforeAll(async () => {
await testDb.init();
repository = Container.get(PollerStateRepository);
workflowRepository = Container.get(WorkflowRepository);
({ id: workflowId } = await createWorkflow());
});
beforeEach(async () => {
await testDb.truncate(['PollerState']);
});
afterAll(async () => {
await testDb.terminate();
});
const seed = async (nodeId: string, cursor: PollerCursor, ofWorkflow = workflowId) =>
// TypeORM's insert type deep-partialises the JSON column into a shape a plain
// `Record<string, unknown>` cannot satisfy.
await repository.insert({
workflowId: ofWorkflow,
nodeId,
cursor: cursor as Record<string, object>,
});
describe('findCursor', () => {
it('returns null for a node that has never polled', async () => {
expect(await repository.findCursor(workflowId, 'node-1')).toBeNull();
});
it('distinguishes an empty cursor from a missing one', async () => {
await seed('node-1', {});
expect(await repository.findCursor(workflowId, 'node-1')).toEqual({});
});
// The cursor column is `json` on Postgres but `text` on SQLite, so the two
// dialects reach the same value by different routes.
it('returns a stored cursor unchanged', async () => {
const cursor = {
lastTimeChecked: '2026-07-28T10:00:00.000Z',
possibleDuplicates: ['id-1', 'id-2'],
};
await seed('node-1', cursor);
expect(await repository.findCursor(workflowId, 'node-1')).toEqual(cursor);
});
});
describe('table constraints', () => {
it('keeps the cursors of two nodes in the same workflow separate', async () => {
await seed('node-1', { lastItemId: 'a' });
await seed('node-2', { lastItemId: 'b' });
expect(await repository.findCursor(workflowId, 'node-1')).toEqual({ lastItemId: 'a' });
expect(await repository.findCursor(workflowId, 'node-2')).toEqual({ lastItemId: 'b' });
});
it('keeps the cursors of one node id in two workflows separate', async () => {
const { id: otherWorkflowId } = await createWorkflow();
await seed('node-1', { lastItemId: 'a' });
await seed('node-1', { lastItemId: 'b' }, otherWorkflowId);
expect(await repository.findCursor(workflowId, 'node-1')).toEqual({ lastItemId: 'a' });
expect(await repository.findCursor(otherWorkflowId, 'node-1')).toEqual({ lastItemId: 'b' });
});
it('rejects a second row for the same workflow and node', async () => {
await seed('node-1', { lastItemId: 'a' });
await expect(seed('node-1', { lastItemId: 'b' })).rejects.toThrow();
expect(await repository.findCursor(workflowId, 'node-1')).toEqual({ lastItemId: 'a' });
});
it('rejects a cursor for a workflow that does not exist', async () => {
await expect(seed('node-1', { lastItemId: 'a' }, 'does-not-exist')).rejects.toThrow();
expect(await repository.findCursor('does-not-exist', 'node-1')).toBeNull();
});
it("drops a workflow's cursors when the workflow is deleted", async () => {
const { id: doomedWorkflowId } = await createWorkflow();
await seed('node-1', { lastItemId: 'a' }, doomedWorkflowId);
await seed('node-2', { lastItemId: 'b' }, doomedWorkflowId);
await workflowRepository.delete({ id: doomedWorkflowId });
expect(await repository.findCursor(doomedWorkflowId, 'node-1')).toBeNull();
expect(await repository.findCursor(doomedWorkflowId, 'node-2')).toBeNull();
});
});
});