diff --git a/airbyte-integrations/connectors/source-hubspot/acceptance-test-config.yml b/airbyte-integrations/connectors/source-hubspot/acceptance-test-config.yml index f0bab1342b56..472d58646a0f 100644 --- a/airbyte-integrations/connectors/source-hubspot/acceptance-test-config.yml +++ b/airbyte-integrations/connectors/source-hubspot/acceptance-test-config.yml @@ -38,6 +38,18 @@ acceptance_tests: bypass_reason: No value in re-create every 90 days - name: list_memberships bypass_reason: Test account may not have populated list memberships + - name: goals + bypass_reason: >- + Test portal lacks Sales Hub Enterprise, so /crm/v3/objects/goal_targets returns 403 + (requires one of [goals-read]) - covered by unit tests + - name: leads + bypass_reason: >- + Test portal lacks Sales Hub Professional, so /crm/v3/properties/leads returns 403 - + covered by unit tests + - name: workflows + bypass_reason: >- + Test portal lacks automation access, so /automation/v3/workflows returns 403 (missing + 'workflows-access-public-api') - covered by unit tests - name: contacts_web_analytics bypass_reason: High-fanout experimental stream validated by unit tests and manual preview - name: companies_web_analytics @@ -73,6 +85,18 @@ acceptance_tests: bypass_reason: No value in re-create every 90 days - name: list_memberships bypass_reason: Test account may not have populated list memberships + - name: goals + bypass_reason: >- + Test portal lacks Sales Hub Enterprise, so /crm/v3/objects/goal_targets returns 403 + (requires one of [goals-read]) - covered by unit tests + - name: leads + bypass_reason: >- + Test portal lacks Sales Hub Professional, so /crm/v3/properties/leads returns 403 - + covered by unit tests + - name: workflows + bypass_reason: >- + Test portal lacks automation access, so /automation/v3/workflows returns 403 (missing + 'workflows-access-public-api') - covered by unit tests - name: contacts_web_analytics bypass_reason: High-fanout experimental stream validated by unit tests and manual preview - name: companies_web_analytics diff --git a/airbyte-integrations/connectors/source-hubspot/manifest.yaml b/airbyte-integrations/connectors/source-hubspot/manifest.yaml index 0dacfe7adba0..56fa8d4fec54 100644 --- a/airbyte-integrations/connectors/source-hubspot/manifest.yaml +++ b/airbyte-integrations/connectors/source-hubspot/manifest.yaml @@ -3356,7 +3356,7 @@ api_budget: # HubSpot account is limited to 110 requests every 10 seconds https://developers.hubspot.com/docs/guides/apps/api-usage/usage-details#rate-limits concurrency_level: type: ConcurrencyLevel - default_concurrency: "{{ config.get('num_workers', 10) }}" + default_concurrency: "{{ config.get('num_worker', 10) }}" max_concurrency: 40 schemas: diff --git a/airbyte-integrations/connectors/source-hubspot/metadata.yaml b/airbyte-integrations/connectors/source-hubspot/metadata.yaml index 4060ee6a199d..def02f4b171a 100644 --- a/airbyte-integrations/connectors/source-hubspot/metadata.yaml +++ b/airbyte-integrations/connectors/source-hubspot/metadata.yaml @@ -10,7 +10,7 @@ data: connectorSubtype: api connectorType: source definitionId: 36c891d9-4bd9-43ac-bad2-10e12756272c - dockerImageTag: 6.8.0 + dockerImageTag: 6.8.1 dockerRepository: airbyte/source-hubspot documentationUrl: https://docs.airbyte.com/integrations/sources/hubspot resourceRequirements: diff --git a/airbyte-integrations/connectors/source-hubspot/unit_tests/test_concurrency.py b/airbyte-integrations/connectors/source-hubspot/unit_tests/test_concurrency.py new file mode 100644 index 000000000000..0442fd769a36 --- /dev/null +++ b/airbyte-integrations/connectors/source-hubspot/unit_tests/test_concurrency.py @@ -0,0 +1,20 @@ +# +# Copyright (c) 2026 Airbyte, Inc., all rights reserved. +# + +import pytest + +from .conftest import get_source + + +@pytest.mark.parametrize( + "config, expected_concurrency", + [ + pytest.param({"num_worker": 25}, 25, id="configured_num_worker"), + pytest.param({}, 10, id="default_num_worker"), + ], +) +def test_concurrency_level_uses_num_worker(config, expected_concurrency): + source = get_source(config) + + assert source._concurrent_source._threadpool._threadpool._max_workers == expected_concurrency diff --git a/airbyte-integrations/connectors/source-hubspot/unit_tests/test_entitlement_gated_streams.py b/airbyte-integrations/connectors/source-hubspot/unit_tests/test_entitlement_gated_streams.py new file mode 100644 index 000000000000..7e938ba37b00 --- /dev/null +++ b/airbyte-integrations/connectors/source-hubspot/unit_tests/test_entitlement_gated_streams.py @@ -0,0 +1,148 @@ +# +# Copyright (c) 2026 Airbyte, Inc., all rights reserved. +# + +"""Coverage for streams bypassed in `acceptance-test-config.yml` via `empty_streams`. + +`goals`, `leads`, and `workflows` back onto HubSpot objects that need paid product access +(Sales Hub Enterprise, Sales Hub Professional, and automation respectively). The CI test portal +lost those entitlements, so the live acceptance read returns 403 and yields no records. These +tests keep the behaviour that `basic_read` used to assert: records come through, they match the +declared schema, incremental runs emit state, and an empty page is not treated as a failure. +""" + +import pytest + +from airbyte_cdk.models import SyncMode +from airbyte_cdk.test.entrypoint_wrapper import discover + +from .conftest import find_stream, get_source, mock_dynamic_schema_requests_with_skip, read_from_stream + + +# stream name -> (properties entity used for dynamic schemas, one API record shaped like the +# stream's declared schema; `workflows` carries an epoch-millis cursor rather than an ISO string) +GATED_STREAMS = { + "goals": ( + "goal_targets", + { + "id": "test_id", + "createdAt": "2022-02-25T16:43:11Z", + "updatedAt": "2022-02-25T16:43:11Z", + "properties": {"hs__test_field": "value"}, + }, + ), + "leads": ( + "leads", + { + "id": "test_id", + "createdAt": "2022-02-25T16:43:11Z", + "updatedAt": "2022-02-25T16:43:11Z", + "properties": {"hs__test_field": "value"}, + }, + ), + "workflows": ( + "", + {"id": "test_id", "insertedAt": 1675121674226, "updatedAt": 1675121674226, "enabled": True}, + ), +} + + +def _retriever(stream_name, config): + stream = find_stream(stream_name, config) + return stream, stream._stream_partition_generator._partition_factory._retriever + + +def _register_stream_response(requests_mock, stream_name, entity, config, records): + """Mock everything a read of `stream_name` touches, returning `records` from its endpoint.""" + mock_dynamic_schema_requests_with_skip(requests_mock, []) + + stream, retriever = _retriever(stream_name, config) + stream._sync_mode = SyncMode.full_refresh + url = retriever.requester.url_base + "/" + retriever.requester.get_path(stream_slice={}) + stream._sync_mode = None + + method = retriever.requester._http_method.value + field_path = retriever.record_selector.extractor.field_path + data_field = field_path[0] if len(field_path) > 0 else None + body = {data_field: records} if data_field else records + requests_mock.register_uri(method, url, [{"json": body, "status_code": 200}]) + + # CRM search streams fan out to association batch reads once they have record ids. + if method == "POST": + for association in retriever.requester._parameters.get("associations", []): + requests_mock.register_uri( + "POST", + f"https://api.hubapi.com/crm/v4/associations/{entity}/{association}/batch/read", + [{"json": {"results": []}, "status_code": 200}], + ) + return stream + + +def _schema_for(stream_name, config, requests_mock): + mock_dynamic_schema_requests_with_skip(requests_mock, []) + catalog = discover(get_source(config), config).catalog.catalog + for stream in catalog.streams: + if stream.name == stream_name: + return stream.json_schema + raise AssertionError(f"{stream_name} missing from the discovered catalog") + + +@pytest.mark.parametrize("stream_name", list(GATED_STREAMS)) +def test_read_emits_records_matching_declared_schema(stream_name, requests_mock, config): + entity, record = GATED_STREAMS[stream_name] + _register_stream_response(requests_mock, stream_name, entity, config, [record]) + + output = read_from_stream(config, stream_name, SyncMode.full_refresh) + + # `leads` fetches properties in chunks, so one API record can surface once per chunk; assert on + # identity rather than a count that would encode chunking internals. + assert output.records + assert {message.record.data["id"] for message in output.records} == {"test_id"} + + schema_properties = _schema_for(stream_name, config, requests_mock)["properties"] + for message in output.records: + unexpected = [field for field in message.record.data if field not in schema_properties] + assert not unexpected, f"{stream_name} emitted fields absent from its schema: {unexpected}" + + +@pytest.mark.parametrize("stream_name", list(GATED_STREAMS)) +def test_incremental_read_emits_state(stream_name, requests_mock, config): + entity, record = GATED_STREAMS[stream_name] + _register_stream_response(requests_mock, stream_name, entity, config, [record]) + + output = read_from_stream(config, stream_name, SyncMode.incremental) + + assert output.records + assert output.state_messages, f"{stream_name} produced no state message" + assert output.most_recent_state.stream_descriptor.name == stream_name + + +@pytest.mark.parametrize("stream_name", list(GATED_STREAMS)) +def test_empty_page_yields_no_records_and_no_error(stream_name, requests_mock, config): + """An entitlement-limited portal legitimately returns nothing; that is not a sync failure.""" + entity, _ = GATED_STREAMS[stream_name] + _register_stream_response(requests_mock, stream_name, entity, config, []) + + output = read_from_stream(config, stream_name, SyncMode.full_refresh) + + assert len(output.records) == 0 + assert len(output.errors) == 0 + + +@pytest.mark.parametrize("stream_name", list(GATED_STREAMS)) +def test_missing_entitlement_403_raises_actionable_error(stream_name, requests_mock, config): + """A portal without the paid product answers 403; the sync must fail loudly, not silently.""" + entity, _ = GATED_STREAMS[stream_name] + mock_dynamic_schema_requests_with_skip(requests_mock, []) + stream, retriever = _retriever(stream_name, config) + stream._sync_mode = SyncMode.full_refresh + url = retriever.requester.url_base + "/" + retriever.requester.get_path(stream_slice={}) + stream._sync_mode = None + requests_mock.register_uri(retriever.requester._http_method.value, url, [{"status_code": 403, "json": {}}]) + + output = read_from_stream(config, stream_name, SyncMode.full_refresh) + + assert len(output.records) == 0 + assert output.errors, f"{stream_name} swallowed a 403" + assert "Access denied (403)" in output.errors[0].trace.error.message + assert f"to access stream {stream_name}" in output.errors[0].trace.error.message diff --git a/docs/integrations/sources/hubspot.md b/docs/integrations/sources/hubspot.md index 3ea350757952..c260d662b1d5 100644 --- a/docs/integrations/sources/hubspot.md +++ b/docs/integrations/sources/hubspot.md @@ -458,6 +458,7 @@ If you use Airbyte Cloud and your organization restricts access to specific IPs, | Version | Date | Pull Request | Subject | |:------------|:-----------|:---------------------------------------------------------|:-----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------| +| 6.8.1 | 2026-08-14 | [84411](https://github.com/airbytehq/airbyte/pull/84411) | Fix the `Number of concurrent threads` setting being ignored: read the `num_worker` config key emitted by the spec instead of `num_workers` | | 6.8.0 | 2026-06-24 | [80806](https://github.com/airbytehq/airbyte/pull/80806) | Restore 12 Web Analytics streams removed during the v5.8.0 manifest-only migration. Uses HubSpot's latest 2026-03 Events API endpoint with explicit `eventType` fanout. Streams require the `business-intelligence` scope (Marketing Hub Enterprise) plus each parent stream's read scope. Gated behind `enable_experimental_streams` with fresh state from `start_date`. | | 6.7.0 | 2026-06-11 | [76396](https://github.com/airbytehq/airbyte/pull/76396) | Add `treat_numbers_and_booleans_as_strings` config toggle to coerce dynamic `number`/`boolean` properties to `string` | | 6.6.1 | 2026-06-10 | [79636](https://github.com/airbytehq/airbyte/pull/79636) | Add configurable `property_history_lookback_window` (minutes) to property history streams (deals, contacts, companies) to prevent silent record loss caused by cursor drift from HubSpot calculated properties. Clarify existing `lookback_window` field as CRM Search-specific. |