Compare commits

...
Author SHA1 Message Date
Victor Tintel b7b7d23ac2 fix: add content scope to HubSpot OAuth 2026-09-02 22:17:57 -03:00
Rafael Santana ec39b82843 Merge pull request #518 from dadosfera/feature/cdc-connector
FIX: CDC live status decodes int64 offsets as numbers
2026-08-31 17:07:20 -03:00
RafaelandWOZCODE 20d3532c0f FIX: address review round 2 on the CDC refactor
- toCdcTable accepts the source table named either table_name (platform
  bodies) or name (create DTO), so call sites pass it point-free:
  body.tables.map(toCdcTable). The identity is still derived in one place.
- PipelineTablesService declares its logger like every other maestro service
  (logger: DadosferaLogger assigned from the injected instance).

Co-Authored-By: WOZCODE <contact@withwoz.com>
2026-08-31 16:09:25 -03:00
Rafael Santana 9457fcc85e Merge pull request #519 from dadosfera/refactor/cdc-review
UPDATE: apply the CDC code-review policy (maestro #510)
2026-08-31 12:52:51 -03:00
Rafael Santana 1400df0ab9 Merge pull request #510 from dadosfera/feature/cdc-connector
Feature/cdc connector
2026-08-28 13:36:54 -03:00
5 changed files with 71 additions and 9 deletions
+9 -4
View File
@@ -3,11 +3,15 @@ import {
CdcTable,
} from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
/** What a caller knows about a CDC table before it is stored. */
/** What a caller knows about a CDC table before it is stored. The source
* table is named `table_name` on the platform-facing bodies and `name` on
* the create DTO (CdcTableReq); either works — deriving one from the other
* happens only here. */
export interface CdcTableInput {
// Optional on the create DTO (CdcTableReq); the platform validates it.
table_schema?: string;
table_name: string;
table_name?: string;
name?: string;
primary_keys?: string[];
iceberg_table_name?: string;
iceberg_qualify_table_name?: string;
@@ -25,10 +29,11 @@ export interface CdcTableInput {
* `name ?? table_name`).
*/
export function toCdcTable(table: CdcTableInput): CdcTable {
const tableName = table.table_name ?? table.name;
return {
table_schema: table.table_schema,
table_name: table.table_name,
name: table.table_name,
table_name: tableName,
name: tableName,
primary_keys: table.primary_keys ?? [],
iceberg_table_name: table.iceberg_table_name,
iceberg_qualify_table_name: table.iceberg_qualify_table_name,
+1 -3
View File
@@ -196,9 +196,7 @@ export class InputsService {
name: body.name,
plugin: body.plugin,
read_only: body.read_only ?? true,
tables: body.tables.map((t) =>
toCdcTable({ ...t, table_name: t.name }),
),
tables: body.tables.map(toCdcTable),
destination: body.destination,
},
info,
@@ -0,0 +1,59 @@
let mockStrategyOptions: { scope: string };
jest.mock('passport-hubspot-oauth2', () => ({
Strategy: class {
constructor(options: { scope: string }) {
mockStrategyOptions = options;
}
authenticate() {}
},
}));
jest.mock('@nestjs/passport', () => ({
PassportStrategy: (Strategy) => Strategy,
}));
import { HubspotStrategy } from './hubspot-strategy';
import type { OauthSecrets } from 'src/utils/OauthSecrets';
import type { OauthService } from '../oauth.service';
describe('HubspotStrategy', () => {
it('requests the content scope without requesting Marketing write scopes', () => {
const oauthSecrets = {
hubspot: {
client_id: '',
client_secret: '',
redirect_uri: '',
},
} as OauthSecrets;
new HubspotStrategy(oauthSecrets, {} as OauthService);
const scopes = mockStrategyOptions.scope.split(' ');
expect(scopes).toEqual([
'tickets',
'automation',
'business-intelligence',
'oauth',
'forms',
'content',
'integration-sync',
'sales-email-read',
'crm.lists.read',
'crm.objects.contacts.read',
'crm.schemas.contacts.read',
'crm.objects.companies.read',
'crm.objects.deals.read',
'crm.schemas.companies.read',
'crm.schemas.deals.read',
'crm.objects.owners.read',
'crm.objects.quotes.read',
'crm.schemas.quotes.read',
'crm.objects.line_items.read',
'crm.schemas.line_items.read',
]);
expect(scopes).not.toContain('marketing.email.write');
});
});
@@ -17,7 +17,7 @@ export class HubspotStrategy extends PassportStrategy(Strategy) {
clientSecret: oauthSecrets.hubspot.client_secret,
callbackURL: oauthSecrets.hubspot.redirect_uri,
scope:
'tickets automation business-intelligence oauth forms integration-sync sales-email-read crm.lists.read crm.objects.contacts.read crm.schemas.contacts.read crm.objects.companies.read crm.objects.deals.read crm.schemas.companies.read crm.schemas.deals.read crm.objects.owners.read crm.objects.quotes.read crm.schemas.quotes.read crm.objects.line_items.read crm.schemas.line_items.read',
'tickets automation business-intelligence oauth forms content integration-sync sales-email-read crm.lists.read crm.objects.contacts.read crm.schemas.contacts.read crm.objects.companies.read crm.objects.deals.read crm.schemas.companies.read crm.schemas.deals.read crm.objects.owners.read crm.objects.quotes.read crm.schemas.quotes.read crm.objects.line_items.read crm.schemas.line_items.read',
passReqToCallback: true,
},
(accessToken, refreshToken, tokenInfo, profile, done) => {
@@ -30,7 +30,7 @@ interface PlatformPipeline {
*/
@Injectable()
export class PipelineTablesService {
private logger: DadosferaLogger['logger'];
logger: DadosferaLogger;
constructor(
private readonly platformApiService: PlatformApiService,