Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions .husky/pre-commit
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
#!/usr/bin/env sh
. "$(dirname -- "$0")/_/husky.sh"


echo "🔍 Running pre-commit checks..."

Expand Down
1 change: 1 addition & 0 deletions apps/api/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
"@nexiom/dbmanager": "workspace:*",
"@nexiom/engine": "workspace:*",
"@nexiom/identity": "workspace:*",
"@nexiom/piece-framework": "workspace:*",
"@nexiom/piece-quickbooks": "workspace:*",
"@nexiom/piece-salesforce": "workspace:*",
"@nexiom/queue": "workspace:*",
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import { Test, TestingModule } from '@nestjs/testing';
import { OAuthCallbackController } from './callback.controller.js';
import { PieceRegistryService } from '@nexiom/engine';
import type { Piece } from '@nexiom/connectors/framework';
import type { Piece } from '@nexiom/piece-framework';
import type { Request, Response } from 'express';
import { vi, describe, it, expect, beforeEach } from 'vitest';
import type { Mocked } from 'vitest';
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import type { Piece } from '@nexiom/connectors/framework';
import type { Piece } from '@nexiom/piece-framework';
/* eslint-disable @typescript-eslint/unbound-method */
import { Test, TestingModule } from '@nestjs/testing';
import { ConnectorsController } from './connectors.controller.js';
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ import type { Response } from 'express';
import { AuthContext, type RequestAuthContext, AuthGuard } from '@nexiom/auth';
import { getAdminRoleId, getOwnerRoleId } from '@nexiom/identity/constants';
import { EncryptionService, AppCredentialError } from '@nexiom/connectors';
import type { AnyProperty } from '@nexiom/connectors';
import type { AnyProperty } from '@nexiom/piece-framework';
import { ConnectorsService } from '../connectors.service.js';
import { OauthStateService } from '../oauth-state.service.js';
import {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { DefaultOAuthRefreshClient } from './token-refresh.service.js';
import { EncryptionService, OAuthRefreshError } from '@nexiom/connectors';
import { PieceRegistryService } from '@nexiom/engine';
import type { Piece } from '@nexiom/connectors/framework';
import type { Piece } from '@nexiom/piece-framework';
import {
describe,
it,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import {
OAuthRefreshError,
EncryptionService,
} from '@nexiom/connectors';
import { PropertyType, resolveOAuth2Url } from '@nexiom/connectors/framework';
import { PropertyType, resolveOAuth2Url } from '@nexiom/piece-framework';
import {
appConnections,
AppConnectionStatus,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ import { ConfigService } from '@nestjs/config';
import { ConnectorsService } from './connectors.service.js';
import { EncryptionService, AppCredentialError } from '@nexiom/connectors';
import { PieceRegistryService } from '@nexiom/engine';
import type { Piece } from '@nexiom/connectors/framework';
import type { Piece } from '@nexiom/piece-framework';
import { DB_MANAGER } from '../dbmanager/dbmanager.module.js';
import { SchemaPlan } from '@nexiom/dbmanager';
import {
Expand Down
10 changes: 4 additions & 6 deletions apps/api/src/modules/connections/connectors.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,10 @@ import {
ConflictException,
} from '@nestjs/common';
import { ConfigService } from '@nestjs/config';
import {
AppCredentialError,
resolveOAuth2Url,
PropertyType,
} from '@nexiom/connectors';
import type { OAuth2Auth, OAuthCredentialBlob } from '@nexiom/connectors';
import { AppCredentialError } from '@nexiom/connectors';
import { resolveOAuth2Url, PropertyType } from '@nexiom/piece-framework';
import type { OAuthCredentialBlob } from '@nexiom/connectors';
import type { OAuth2Auth } from '@nexiom/piece-framework';
import {
appConnections,
AppConnectionStatus,
Expand Down
3 changes: 2 additions & 1 deletion apps/api/src/modules/scheduler/poll-sync-runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,8 @@ import {
syncCursors,
} from '@nexiom/database';
import { TokenManagerService } from '@nexiom/connectors';
import type { Piece, OAuthCredentialBlob } from '@nexiom/connectors';
import type { OAuthCredentialBlob } from '@nexiom/connectors';
import type { Piece } from '@nexiom/piece-framework';
import { REDIS_CLIENT, type Redis } from '@nexiom/cache';
import {
CursorManagerService,
Expand Down
4 changes: 2 additions & 2 deletions apps/api/src/modules/stitches/metadata-discovery.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,11 @@ import {
import { REDIS_CLIENT, type Redis } from '@nexiom/cache';
import { TokenManagerService } from '@nexiom/connectors';
import { PieceRegistryService } from '@nexiom/engine';
import type { OAuthCredentialBlob } from '@nexiom/connectors';
import type {
ObjectDescriptor,
FieldDescriptor,
OAuthCredentialBlob,
} from '@nexiom/connectors';
} from '@nexiom/piece-framework';

// Single source of truth for metadata cache TTL.
const TTL_SECONDS = 5 * 60; // 5 minutes
Expand Down
2 changes: 1 addition & 1 deletion apps/api/src/modules/trigger/dlq-processor.service.spec.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { describe, it, expect, beforeEach, afterEach, vi } from 'vitest';
import { DlqProcessorService } from './dlq-processor.service.js';
/* eslint-disable @typescript-eslint/unbound-method */
import { TriggerStrategy } from '@nexiom/connectors';
import { TriggerStrategy } from '@nexiom/piece-framework';
import type { TriggerExecutorService } from './trigger-executor.service.js';
import { PieceRegistryService } from '@nexiom/engine';

Expand Down
2 changes: 1 addition & 1 deletion apps/api/src/modules/trigger/poller.service.spec.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { describe, it, expect, beforeEach, vi } from 'vitest';
import { PollerService } from './poller.service.js';
/* eslint-disable @typescript-eslint/unbound-method */
import { TriggerStrategy } from '@nexiom/connectors';
import { TriggerStrategy } from '@nexiom/piece-framework';
import type { TriggerExecutorService } from './trigger-executor.service.js';
import { PieceRegistryService } from '@nexiom/engine';

Expand Down
2 changes: 1 addition & 1 deletion apps/api/src/modules/trigger/redis-trigger-store.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import type { TriggerStore } from '@nexiom/connectors';
import type { TriggerStore } from '@nexiom/piece-framework';
import type { Redis } from 'ioredis';

/**
Expand Down
4 changes: 2 additions & 2 deletions apps/api/src/modules/trigger/trigger-executor.service.spec.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import { describe, it, expect, beforeEach, vi } from 'vitest';
import { TriggerExecutorService } from './trigger-executor.service.js';
import { TriggerStrategy } from '@nexiom/connectors';
import type { Trigger } from '@nexiom/connectors';
import { TriggerStrategy } from '@nexiom/piece-framework';
import type { Trigger } from '@nexiom/piece-framework';

function makeMockDb() {
return {
Expand Down
2 changes: 1 addition & 1 deletion apps/api/src/modules/trigger/trigger-executor.service.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import { Injectable, Inject, Logger } from '@nestjs/common';
import type { Trigger, TriggerContext } from '@nexiom/connectors';
import type { Trigger, TriggerContext } from '@nexiom/piece-framework';
import type { DrizzleDb } from '@nexiom/database';
import { DATABASE_CONNECTION } from '@nexiom/database';
import { SchemaPlan } from '@nexiom/dbmanager';
Expand Down
2 changes: 1 addition & 1 deletion apps/api/src/modules/trigger/webhooks.controller.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
import { describe, it, expect, beforeEach, vi, type Mock } from 'vitest';
import { WebhooksController } from './webhooks.controller.js';
import { NotFoundException, UnauthorizedException } from '@nestjs/common';
import { TriggerStrategy } from '@nexiom/connectors';
import { TriggerStrategy } from '@nexiom/piece-framework';
import { PieceRegistryService } from '@nexiom/engine';
import type { TriggerExecutorService } from './trigger-executor.service.js';

Expand Down
1 change: 1 addition & 0 deletions apps/worker/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@
"@nexiom/dbmanager": "workspace:*",
"@nexiom/engine": "workspace:*",
"@nexiom/identity": "workspace:*",
"@nexiom/piece-framework": "workspace:*",
"@nexiom/piece-quickbooks": "workspace:*",
"@nexiom/piece-salesforce": "workspace:*",
"@nexiom/queue": "workspace:*",
Expand Down
2 changes: 1 addition & 1 deletion apps/worker/src/modules/pipeline/delivery.service.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -390,7 +390,7 @@ describe("DeliveryService", () => {
});

it("should set RETRY status when piece throws RetryableException", async () => {
const { RetryableException } = await import("@nexiom/connectors");
const { RetryableException } = await import("@nexiom/piece-framework");
pieceRegistry
.getPiece()
.executeAction.mockRejectedValueOnce(
Expand Down
5 changes: 2 additions & 3 deletions apps/worker/src/modules/pipeline/delivery.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import {
Optional,
Logger,
} from "@nestjs/common";
import { eq } from "drizzle-orm";
import { eq, sql } from "drizzle-orm";
import { QueueService, QueueName } from "@nexiom/queue";
import {
DATABASE_CONNECTION,
Expand All @@ -18,9 +18,8 @@ import {
} from "@nexiom/database";
import type { DrizzleDb } from "@nexiom/database";
import { StorageResolverService, PieceRegistryService } from "@nexiom/engine";
import { sql } from "drizzle-orm";
import { TokenManagerService } from "@nexiom/connectors";
import { RetryableException } from "@nexiom/connectors";
import { RetryableException } from "@nexiom/piece-framework";
import {
sanitizeError,
isValidPipelineMessage,
Expand Down
3 changes: 2 additions & 1 deletion apps/worker/src/modules/pipeline/fanout.service.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,8 @@ describe("FanOutService", () => {
}),
onConflictDoNothing: vi.fn().mockResolvedValue(undefined),
returning: vi.fn().mockResolvedValue([{ id: "outbound_1" }]),
then: (res: any) => Promise.resolve(undefined).then(res),
then: (onfulfilled?: ((value: any) => any) | null) =>
Promise.resolve(undefined as any).then(onfulfilled),
}),
});
queueService = { consume: vi.fn(), send: vi.fn() };
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,15 @@ describe("NormalizationService", () => {
values: vi.fn().mockReturnValue({
onConflictDoNothing: vi.fn().mockReturnValue({
// directly awaitable (for normalizedOutbox insert that has no .returning())
then: (res: any) => Promise.resolve(undefined).then(res),
then: (onfulfilled?: ((value: any) => any) | null) =>
Promise.resolve(undefined as any).then(onfulfilled),
// also supports .returning() for chains that need it
returning: vi.fn().mockResolvedValue(returnVal),
}),
returning: vi.fn().mockResolvedValue(returnVal),
// plain insert().values() with no conflict resolution
then: (res: any) => Promise.resolve(undefined).then(res),
then: (onfulfilled?: ((value: any) => any) | null) =>
Promise.resolve(undefined as any).then(onfulfilled),
}),
});
mockTxInsert.mockImplementation(() => makeInsertChain());
Expand Down Expand Up @@ -158,7 +160,8 @@ describe("NormalizationService", () => {
onConflictDoNothing: vi.fn().mockReturnValue({
returning: vi.fn().mockResolvedValue([]),
}),
then: (res: any) => Promise.resolve(undefined).then(res),
then: (onfulfilled?: ((value: any) => any) | null) =>
Promise.resolve(undefined as any).then(onfulfilled),
}),
}),
update: vi.fn().mockReturnThis(),
Expand Down
3 changes: 2 additions & 1 deletion apps/worker/src/modules/pipeline/replica.service.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,8 @@ describe("ReplicaService", () => {
returning: vi.fn().mockResolvedValue([{ id: "1" }]),
}),
onConflictDoNothing: vi.fn().mockResolvedValue(undefined),
then: (res: any) => Promise.resolve(undefined).then(res),
then: (onfulfilled?: ((value: any) => any) | null) =>
Promise.resolve(undefined as any).then(onfulfilled),
};
}),
}));
Expand Down
22 changes: 21 additions & 1 deletion docs/architecture/master/tasks.md
Original file line number Diff line number Diff line change
Expand Up @@ -612,6 +612,25 @@ Each task is one commit (or one small PR). Checkboxes track completion.

---

## Phase 8 — Fleet Sharding

### T053 · engine: Custom Logic Extension Hook (`isolated-vm`)

- [ ] Add `isolated-vm` dependency to `@nexiom/engine`.
- [ ] Implement `LogicResolverService` that checks the local shard directory for a tenant's custom script before falling back to generic pipeline logic.
- [ ] Run custom scripts within a secure `ivm.Isolate` context with strict memory buffers and timeouts (e.g., 128MB, 1s timeout) to prevent platform DoS.
- Files: `packages/engine/src/executor/logic-resolver.ts`
- Depends: T035

### T054 · worker: GitOps Shard Synchronization

- [ ] Implement a worker cron service that pulls/fetches mapped Git Shard repositories (e.g. `fluxnex-shard-001`) onto the local disk every 5 minutes.
- [ ] Implement an in-memory or Redis-backed cache invalidation when a shard is updated so the `LogicResolverService` uses the latest custom logic from customers.
- Files: `apps/worker/src/modules/gitops/shard-sync.service.ts`
- Depends: T053

---

## Summary

| Phase | Tasks | Completed | Key deliverable |
Expand All @@ -626,8 +645,9 @@ Each task is one commit (or one small PR). Checkboxes track completion.
| 5 — AI Mapping | T040–T042 | ⬜ All | Claude-powered field suggestions |
| 6 — Environments | T043–T045 | ⬜ All | Sandbox/Production routing |
| 7 — Delivery Outbox | T051–T052 | ✅ All | Delivery Outbox Resiliency |
| 8 — Fleet Sharding | T053–T054 | ⬜ All | Sandboxed execution of customer logic |

**Total: 52 tasks · Completed: ~38 · Remaining: ~14**
**Total: 54 tasks · Completed: ~38 · Remaining: ~16**

---

Expand Down
6 changes: 1 addition & 5 deletions packages/connectors/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -11,11 +11,6 @@
"import": "./dist/index.js",
"default": "./dist/index.js"
},
"./framework": {
"types": "./dist/framework/index.d.ts",
"import": "./dist/framework/index.js",
"default": "./dist/framework/index.js"
},
"./intelligence": {
"types": "./dist/intelligence/index.d.ts",
"import": "./dist/intelligence/index.js",
Expand All @@ -30,6 +25,7 @@
},
"dependencies": {
"@nexiom/database": "workspace:*",
"@nexiom/piece-framework": "workspace:*",
"drizzle-orm": "^0.45.1",
"ioredis": "^5.3.2"
},
Expand Down
Loading
Loading