Understanding the plugin architecture helps you use it effectively and extend it.
┌─────────────────────────────────────────────────────────────────┐
│ Vendure Dashboard │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────────────┐ │
│ │ Pipeline │ │ Connections │ │ Runs / Logs / │ │
│ │ Builder │ │ & Secrets │ │ Analytics │ │
│ └──────────────┘ └──────────────┘ └──────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
│
│ GraphQL
▼
┌─────────────────────────────────────────────────────────────────┐
│ Vendure Server │
│ │
│ ┌────────────────────────────────────────────────────────────┐ │
│ │ DataHub Plugin │ │
│ │ │ │
│ │ ┌─────────────┐ ┌─────────────┐ ┌─────────────────────┐ │ │
│ │ │ GraphQL │ │ Job Queue │ │ Webhook │ │ │
│ │ │ Resolvers │ │ Handlers │ │ Controllers │ │ │
│ │ └─────────────┘ └─────────────┘ └─────────────────────┘ │ │
│ │ │ │ │ │ │
│ │ ▼ ▼ ▼ │ │
│ │ ┌────────────────────────────────────────────────────────┐│ │
│ │ │ Pipeline Runner Service ││ │
│ │ │ ││ │
│ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ ││ │
│ │ │ │ Extract │ │Transform │ │ Validate │ │ Load │ ││ │
│ │ │ │ Executor │ │ Executor │ │ Executor │ │Executor│ ││ │
│ │ │ └──────────┘ └──────────┘ └──────────┘ └────────┘ ││ │
│ │ │ │ │ │ │ ││ │
│ │ │ ▼ ▼ ▼ ▼ ││ │
│ │ │ ┌──────────────────────────────────────────────────┐ ││ │
│ │ │ │ Adapter Registry │ ││ │
│ │ │ │ ┌──────────┐ ┌──────────┐ ┌─────────────────┐ │ ││ │
│ │ │ │ │Extractors│ │Operators │ │ Loaders │ │ ││ │
│ │ │ │ └──────────┘ └──────────┘ └─────────────────┘ │ ││ │
│ │ │ └──────────────────────────────────────────────────┘ ││ │
│ │ └────────────────────────────────────────────────────────┘│ │
│ │ │ │
│ │ ┌────────────────────────────────────────────────────────┐│ │
│ │ │ Services ││ │
│ │ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────┐ ││ │
│ │ │ │ Pipeline │ │ Secrets │ │Connection│ │ Logging│ ││ │
│ │ │ └──────────┘ └──────────┘ └──────────┘ └────────┘ ││ │
│ │ └────────────────────────────────────────────────────────┘│ │
│ └─────────────────────────────────────────────────────────────┘│
│ │
│ ┌────────────────────────────────────────────────────────────┐ │
│ │ Database (TypeORM) │ │
│ │ ┌──────────┐ ┌──────────┐ ┌──────────┐ ┌────────────┐ │ │
│ │ │ Pipeline │ │ Runs │ │ Secrets │ │ Connections│ │ │
│ │ └──────────┘ └──────────┘ └──────────┘ └────────────┘ │ │
│ └────────────────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
DataHubPlugin is the main plugin class that registers:
| Entity | Purpose |
|---|---|
Pipeline |
Pipeline definitions |
PipelineRun |
Execution history, published revision, immutable definition snapshot, initiating channel, and durable queue request |
PipelineRevision |
Version history |
PipelineLog |
Execution logs |
DataHubConnection |
External connections |
DataHubSecret |
Encrypted credentials |
DataHubSettings |
Plugin configuration |
DataHubRecordError |
Failed records |
DataHubCheckpoint |
Resume checkpoints |
DataHubRecordRetryAudit |
Retry audit trail |
| Service | Responsibility |
|---|---|
PipelineService |
CRUD operations for pipelines |
PipelineRunnerService |
Restores the run’s persisted channel and orchestrates execution |
AdapterRuntimeService |
Executes adapters |
SecretService |
Manages secrets |
ConnectionService |
Manages connections |
CheckpointService |
Manages checkpoints |
PipelineLogService |
Writes execution logs |
AnalyticsService |
Aggregates metrics |
RecordErrorService |
Manages failed records |
Step-type-specific execution logic:
| Executor | Step Types |
|---|---|
ExtractExecutor |
extract |
TransformExecutor |
transform, validate, enrich |
LoadExecutor |
load |
ExportExecutor |
export |
FeedExecutor |
feed |
SinkExecutor |
sink |
The registry manages all adapters via registerRuntime():
// All adapter types use the same registration method
registry.registerRuntime(httpApiExtractor);
registry.registerRuntime(databaseExtractor);
registry.registerRuntime(renameOperator);
registry.registerRuntime(setOperator);
registry.registerRuntime(productLoader);
registry.registerRuntime(customerLoader);
A pipeline run starts when:
Runs are processed via Vendure’s data-hub.run queue. Creating or resuming a
run stores a durable queue request on PipelineRun in the same database
transaction as the lifecycle change. The queue handler atomically claims that
request before adding the Vendure job and periodically recovers missing or
stale claims. The runner clears the request only after it holds the distributed
execution lock.
The runner orchestrates execution:
For each step:
Records flow through steps:
Extract → [record1, record2, ...] → Transform → Load
Each step can:
{
version: 1,
steps: [
{ key: 'extract', type: 'extract', config: {...} },
{ key: 'transform', type: 'transform', config: {...} },
{ key: 'load', type: 'load', config: {...} },
],
edges: [
{ from: 'extract', to: 'transform' },
{ from: 'transform', to: 'load' },
],
}
Each run has a context:
{
runId: string;
pipelineId: string;
startedAt: Date;
triggeredBy: string;
parameters: Record<string, any>;
variables: Record<string, any>;
checkpoint: CheckpointData;
}
Records are JSON objects:
interface Record {
[key: string]: JsonValue;
_meta?: {
sourceStep: string;
index: number;
hash?: string;
};
}
Implement the ExtractorAdapter interface:
interface ExtractorAdapter {
readonly type: 'EXTRACTOR';
code: string;
name: string;
extract(context: ExtractContext, config: JsonObject): AsyncGenerator<RecordEnvelope>;
}
Create operator definitions:
interface SingleRecordOperator<TConfig = JsonObject> {
readonly type: 'OPERATOR';
readonly pure: boolean;
applyOne(record: JsonObject, config: TConfig, helpers: AdapterOperatorHelpers): JsonObject | null;
}
Implement entity loading:
interface LoaderAdapter<TConfig = JsonObject> {
readonly type: 'LOADER';
load(context: LoadContext, config: TConfig, records: readonly JsonObject[]): Promise<LoadResult>;
}
See Extending the Plugin for details.
Code-first configuration has two startup paths:
configurationSource: CODE_FIRST. Unchanged definitions are not rewritten, but a matching database-owned row is adopted explicitly. Pipeline changes use the normal lifecycle service and remain non-executable until they pass the review and publish workflow.DATABASE ownership without deleting their rows, revision history, or runs. Dashboard and API mutations reject active CODE_FIRST resources; review and publication remain available for managed pipelines.Custom permissions protect operations:
@Allow(DataHubPipelinePermission.Read)
@Query()
dataHubPipelines() { ... }
@Allow(RunDataHubPipelinePermission.Permission)
@Mutation()
startDataHubPipelineRun() { ... }
Database INLINE secrets require DATAHUB_MASTER_KEY with at least 32 characters and use AES-256-GCM envelopes. Without a valid key, database INLINE writes and resolution fail closed. Unencrypted database values are never resolved.
Code-first INLINE secrets are different: their configuration value is already plaintext in TypeScript, JSON, or YAML. A master key cannot protect that source, so production startup rejects code-first INLINE definitions. Environment-backed definitions accept one canonical variable name and resolve the environment separately in each executing API server or worker.
Webhook requests can be verified with HMAC signatures.
Webhook observation hooks persist a channel-scoped delivery row before queueing network work. The row is the recovery authority; Vendure jobs contain only its database ID and a short-lived lease token. Expired leases and failed queue publications return to the dispatcher without deleting pending work.
The URL, payload, ordinary headers, Secret Code references, and retry policy are stored together in an AES-256-GCM replay envelope. Signing and sensitive header values are resolved from Secret Codes only inside each worker attempt. GraphQL returns a sanitized URL, payload hash and size, status, attempts, timestamps, HTTP status, and bounded safe error text; it never returns the envelope, request headers, payload, secrets, tokens, or response bodies.
Delivery is at-least-once. A channel-scoped idempotency key returns the existing delivery only when its webhook ID and payload fingerprint match; conflicting reuse fails. Receivers must also enforce the transmitted idempotency key.
Records are processed in batches:
{
throughput: {
batchSize: 100,
concurrency: 4,
}
}
Checkpoint persistence is adapter-specific. File, database, CDC, export, file watch, and gate components store only the offsets, cursors, or approval state they explicitly implement. Dirty execution state is normally persisted when a run finalizes; there is no generic periodic last-successful-record scheduler.
Pipeline runs use the job queue:
The worker reloads the selected user and current roles before creating the
Vendure RequestContext. Database-managed runs without a durable actor fail
closed; the code-first superadmin fallback is not applied to them.
Pipelines execute as Directed Acyclic Graphs (DAGs):
┌─────────┐ ┌───────────┐ ┌─────────┐
│ Trigger │────▶│ Extract │────▶│Transform│
└─────────┘ └───────────┘ └────┬────┘
│
┌──────────────────┴──────────────────┐
▼ ▼
┌───────────┐ ┌───────────┐
│ Route │ │ Enrich │
└─────┬─────┘ └─────┬─────┘
┌───────┴───────┐ │
▼ ▼ ▼
┌────────┐ ┌────────┐ ┌───────────┐
│ Load A │ │ Load B │ │ Sink │
└────────┘ └────────┘ └───────────┘
Features:
For multi-instance deployments:
┌─────────────────────────────────────────────────────────────┐
│ Lock Backend Selection │
├─────────────────┬─────────────────┬─────────────────────────┤
│ Redis │ PostgreSQL │ Memory │
│ (Recommended) │ (Fallback) │ (Single Instance) │
├─────────────────┼─────────────────┼─────────────────────────┤
│ SET NX EX │ Advisory Locks │ Map<string, LockToken> │
│ Atomic ops │ pg_try_advisory │ setTimeout cleanup │
│ TTL expiration │ Transaction- │ No cluster support │
│ Cluster support │ scoped │ │
└─────────────────┴─────────────────┴─────────────────────────┘
Configuration:
DataHubPlugin.init({
distributedLock: {
backend: 'redis',
redis: { host: 'localhost', port: 6379 },
defaultTtlMs: 30000,
waitTimeoutMs: 5000,
},
})
Protects external service calls:
┌─────────────────────────────────────────────┐
│ Circuit States │
├─────────────────────────────────────────────┤
│ │
│ ┌────────┐ 5 failures ┌────────┐ │
│ │ CLOSED │─────────────▶│ OPEN │ │
│ │(normal)│ │(blocked)│ │
│ └────┬───┘ └────┬───┘ │
│ ▲ │ │
│ │ ┌───────────┐ │ 30s │
│ 3 success │HALF-OPEN │◀─────┘ │
│ └────│ (testing) │ │
│ └───────────┘ │
└─────────────────────────────────────────────┘
Applied to:
Message queue integration for event-driven pipelines:
┌─────────────────────────────────────────────────────────────┐
│ Queue Adapters │
├──────────────┬──────────────┬─────────────┬────────────────┤
│ RabbitMQ │ Amazon SQS │ Redis │ Internal │
│ (AMQP) │ │ Streams │ (in-process) │
├──────────────┼──────────────┼─────────────┼────────────────┤
│ Native AMQP │ AWS SDK │ XREAD/XADD │ Redis-backed │
│ confirms │ Visibility │ groups │ buffer │
│ basic.consume│ timeout │ Pending │ Manual ack │
│ Prefetch │ Batch recv │ entries │ Dev/test only │
└──────────────┴──────────────┴─────────────┴────────────────┘
Consumer patterns:
Pipelines support multiple concurrent triggers:
┌─────────────────────────────────────────────────────────────┐
│ Trigger Types │
├──────────┬──────────┬──────────┬──────────┬────────────────┤
│ Manual │ Schedule │ Webhook │ Event │ Message Queue │
├──────────┼──────────┼──────────┼──────────┼────────────────┤
│ UI/API │ Cron │ HTTP │ Vendure │ RabbitMQ/SQS/ │
│ trigger │ express │ endpoint │ events │ Redis Streams │
└──────────┴──────────┴──────────┴──────────┴────────────────┘
│ │ │ │ │
└───────────┴──────────┴──────────┴───────────┘
│
▼
┌─────────────────┐
│ Pipeline Runner │
└─────────────────┘
All triggers converge to the same execution engine, enabling:
Product feed generation for marketing channels:
┌─────────────────────────────────────────────────────────────┐
│ Feed Generation Pipeline │
│ │
│ ┌────────────┐ ┌────────────┐ ┌──────────────────┐ │
│ │ Vendure │───▶│ Feed │───▶│ Format Generator │ │
│ │ Products │ │ Filters │ │ │ │
│ └────────────┘ └────────────┘ │ ┌──────────────┐ │ │
│ │ │ Google XML │ │ │
│ │ ├──────────────┤ │ │
│ │ │ Meta/FB │ │ │
│ │ ├──────────────┤ │ │
│ │ │ RSS/Atom │ │ │
│ │ ├──────────────┤ │ │
│ │ │ CSV/JSON │ │ │
│ │ └──────────────┘ │ │
│ └──────────────────┘ │
└─────────────────────────────────────────────────────────────┘
Supports:
Output to external systems:
┌─────────────────────────────────────────────────────────────┐
│ Sink Executor │
│ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ Circuit Breaker │ │
│ └──────────────────────────────────────────────────────┘ │
│ │ │
│ ┌───────────────┼───────────────┐ │
│ ▼ ▼ ▼ │
│ ┌───────────┐ ┌───────────┐ ┌───────────┐ │
│ │ Search │ │ Webhook │ │ Queue │ │
│ │ Engines │ │ │ │ Producer │ │
│ ├───────────┤ ├───────────┤ ├───────────┤ │
│ │MeiliSearch│ │ HTTP POST │ │ RabbitMQ │ │
│ │Elastic │ │ Auth │ │ SQS │ │
│ │Algolia │ │ Retry │ │ Redis │ │
│ │Typesense │ │ Timeout │ │ │ │
│ └───────────┘ └───────────┘ └───────────┘ │
└─────────────────────────────────────────────────────────────┘
All sinks feature:
src/plugins/data-hub/
├── src/ # Backend source
│ ├── api/ # GraphQL resolvers & schema
│ ├── bootstrap/ # Plugin initialization
│ ├── constants/ # Configuration constants
│ ├── decorators/ # Custom decorators
│ ├── enrichers/ # Record enrichers
│ ├── entities/ # TypeORM entities
│ ├── extractors/ # Data extractors
│ ├── feeds/ # Feed generators
│ ├── gql/ # Generated GraphQL types
│ ├── jobs/ # Job queue handlers
│ ├── loaders/ # Entity loaders
│ ├── mappers/ # Field mappers
│ ├── operators/ # Transform operators
│ ├── parsers/ # File parsers
│ ├── runtime/ # Execution engine
│ ├── sdk/ # Public SDK & DSL
│ ├── services/ # Business logic
│ ├── templates/ # Import/export templates
│ ├── transforms/ # Transform execution
│ ├── types/ # TypeScript types
│ ├── utils/ # Utilities
│ ├── validation/ # Pipeline definition validators
│ └── vendure-schemas/ # Vendure entity schema definitions
├── connectors/ # External system connectors (e.g. Pimcore)
├── dashboard/ # Vendure Dashboard extension
│ ├── components/ # UI components
│ ├── constants/ # UI constants
│ ├── gql/ # GraphQL queries
│ ├── hooks/ # React hooks
│ ├── routes/ # Route components
│ ├── types/ # UI type definitions
│ └── utils/ # UI utilities
├── shared/ # Shared code (backend + dashboard)
│ ├── constants/ # Shared constants
│ ├── types/ # Shared types
│ └── utils/ # Shared utilities
├── dev-server/ # Development server & examples
├── docs/ # Documentation
└── e2e/ # End-to-end tests