diff --git a/.github/workflows/10_feature_dbt_checks.yml b/.github/workflows/10_feature_dbt_checks.yml index 90e040068..f25d28e51 100644 --- a/.github/workflows/10_feature_dbt_checks.yml +++ b/.github/workflows/10_feature_dbt_checks.yml @@ -124,12 +124,12 @@ jobs: # There is an issue with --empty and dynamic tables so need to exclude them - name: Governance run of dbt with EMPTY models using slim mode if: ${{ steps.prod_manifest.outputs.manifest_found == 'true' }} - run: "dbt build --fail-fast --defer --state logs --select state:modified+ --empty --exclude config.materialized:dynamic_table ${{ env.FULL_REFRESH_FLAG }}" + run: "dbt build --fail-fast --defer --state logs --select state:modified+ --empty --exclude config.materialized:dynamic_table tag:requires_fixture_data ${{ env.FULL_REFRESH_FLAG }}" # There is an issue with --empty and dynamic tables so need to exclude - name: Governance run of dbt with EMPTY models using full run if: ${{ steps.prod_manifest.outputs.manifest_found == 'false' }} - run: "dbt build --fail-fast --empty --exclude config.materialized:dynamic_table ${{ env.FULL_REFRESH_FLAG }}" + run: "dbt build --fail-fast --empty --exclude config.materialized:dynamic_table tag:requires_fixture_data ${{ env.FULL_REFRESH_FLAG }}" - name: Generate Docs Combining Prod and branch catalog.json if: ${{ steps.prod_manifest.outputs.catalog_found == 'true' }} diff --git a/orchestrate/dags/daily_loan_run__dc_test_cluster.py b/orchestrate/dags/daily_loan_run__dc_test_cluster.py new file mode 100644 index 000000000..5c53aa293 --- /dev/null +++ b/orchestrate/dags/daily_loan_run__dc_test_cluster.py @@ -0,0 +1,100 @@ +""" +## Sample DAG showing end-to-end ELT (DC test cluster) +This DAG shows how to load data with 3 tools, then run dbt, then other tasks. + +Same pipeline as `daily_loan_run`, kept separate because the Airbyte +connection id differs between Datacoves clusters. +""" + +from airflow.decorators import dag, task, task_group +from orchestrate.utils import datacoves_utils + +from fivetran_provider_async.operators import FivetranOperator +from fivetran_provider_async.sensors import FivetranSensor +from airflow.providers.airbyte.operators.airbyte import AirbyteTriggerSyncOperator + +@dag( + doc_md = __doc__, + catchup = False, + + default_args = datacoves_utils.set_default_args( + owner = "Noel Gomez", + owner_email = "noel@example.com" + ), + + description = "Sample DAG to synchronize the Airflow database (DC test cluster)", + schedule = datacoves_utils.set_schedule("0 0 1 */12 *"), + tags=["extract_and_load", "transform", "marketing_automation", "update_catalog"], +) +def daily_loan_run__dc_test_cluster(): + + @task_group( + group_id="extract_and_load_airbyte", + tooltip="Airbyte Extract and Load" + ) + def extract_and_load_airbyte(): + # Extact and load + sync_airbyte = AirbyteTriggerSyncOperator( + task_id="country_populations_datacoves_snowflake", + connection_id="9cccf1c4-978e-4b0a-9677-43c0111dcda9", + airbyte_conn_id="airbyte_connection", + ) + + + @task_group( + group_id="extract_and_load_fivetran", + tooltip="Fivetran Extract and Load" + ) + def extract_and_load_fivetran(): + trigger_fivetran = FivetranOperator( + task_id="datacoves_snowflake_google_analytics_4_trigger", + fivetran_conn_id="fivetran_connection", + connector_id="speak_menial", + wait_for_completion=False, + ) + sensor_fivetran = FivetranSensor( + task_id="datacoves_snowflake_google_analytics_4_sensor", + fivetran_conn_id="fivetran_connection", + connector_id="speak_menial", + poke_interval=60, + ) + trigger_fivetran >> sensor_fivetran + + + @task_group( + group_id="extract_and_load_dlt", + tooltip="dlt Extract and Load" + ) + def extract_and_load_dlt(): + @task.datacoves_bash( + env = datacoves_utils.set_dlt_env_vars({"destinations": ["main_load_keypair"]}), + append_env=True + ) + def load_loans_data(): + + return "cd load/dlt && ./loans_data.py" + + load_loans_data() + + + # Transform Data + @task.datacoves_dbt( + connection_id="main_key_pair" + ) + def transform(): + return "dbt build -s 'tag:daily_run_airbyte+ tag:daily_run_fivetran+'" + + + # Post transformation tasks + @task.datacoves_bash + def marketing_automation(): + return "echo 'send data to marketing tool'" + + @task.datacoves_bash + def update_catalog(): + return "echo 'refresh data catalog'" + + extract_and_load = [extract_and_load_airbyte(), extract_and_load_fivetran(), extract_and_load_dlt()] + extract_and_load >> transform() >> [marketing_automation(), update_catalog()] + +daily_loan_run__dc_test_cluster() diff --git a/orchestrate/dags/notifications_examples/yaml_slack_dag.py b/orchestrate/dags/notifications_examples/yaml_slack_dag.py index 9de28366e..55b3d1017 100644 --- a/orchestrate/dags/notifications_examples/yaml_slack_dag.py +++ b/orchestrate/dags/notifications_examples/yaml_slack_dag.py @@ -1,8 +1,13 @@ import datetime + from airflow.decorators import dag, task from airflow.providers.slack.notifications.slack import send_slack_notification + @dag( + description="Sample DAG with Slack notification, custom image, and resource requests", + schedule="0 0 1 */12 *", + tags=["transform", "slack_notification"], default_args={ "start_date": datetime.datetime(2024, 1, 1, 0, 0), "owner": "Noel Gomez", @@ -10,9 +15,6 @@ "email_on_failure": True, "retries": 3, }, - description="Sample DAG with Slack notification, custom image, and resource requests", - schedule="0 0 1 */12 *", - tags=["transform", "slack_notification"], catchup=False, on_success_callback=send_slack_notification( text="The DAG {{ dag.dag_id }} succeeded", channel="#general" @@ -22,12 +24,13 @@ ), ) def yaml_slack_dag(): - - @task.datacoves_dbt(connection_id="main_key_pair") + @task.datacoves_dbt( + connection_id="main_key_pair", + ) def transform(): return "dbt run -s personal_loans" - transform() + transform = transform() + -# Invoke DAG dag = yaml_slack_dag() diff --git a/orchestrate/dags/notifications_examples/yaml_teams_dag.py b/orchestrate/dags/notifications_examples/yaml_teams_dag.py index 581d449c0..365fac2a3 100644 --- a/orchestrate/dags/notifications_examples/yaml_teams_dag.py +++ b/orchestrate/dags/notifications_examples/yaml_teams_dag.py @@ -1,8 +1,13 @@ import datetime + from airflow.decorators import dag, task from notifiers.datacoves.ms_teams import MSTeamsNotifier + @dag( + description="Sample DAG with MS Teams notification", + schedule="0 0 1 */12 *", + tags=["transform", "ms_teams_notification"], default_args={ "start_date": datetime.datetime(2024, 1, 1, 0, 0), "owner": "Noel Gomez", @@ -10,9 +15,6 @@ "email_on_failure": True, "retries": 3, }, - description="Sample DAG with MS Teams notification", - schedule="0 0 1 */12 *", - tags=["transform", "ms_teams_notification"], catchup=False, on_success_callback=MSTeamsNotifier( connection_id="DATACOVES_MS_TEAMS", theme_color="0000FF" @@ -22,12 +24,13 @@ ), ) def yaml_teams_dag(): - - @task.datacoves_dbt(connection_id="main_key_pair") + @task.datacoves_dbt( + connection_id="main_key_pair", + ) def transform(): return "dbt run -s personal_loans" - transform() + transform = transform() + -# Invoke DAG dag = yaml_teams_dag() diff --git a/orchestrate/dags/other_examples/dbt_templated_overrides_dag.py b/orchestrate/dags/other_examples/dbt_templated_overrides_dag.py new file mode 100644 index 000000000..e14493823 --- /dev/null +++ b/orchestrate/dags/other_examples/dbt_templated_overrides_dag.py @@ -0,0 +1,50 @@ +""" +## Sample DAG showing templated overrides on the dbt decorator +This DAG shows how to parameterize the `@task.datacoves_dbt` `overrides` kwarg +at trigger time using Jinja templates and DAG run `params`, instead of having +to hardcode the schema/role/warehouse or fall back to an Airflow Variable. +""" + +from airflow.decorators import dag, task +from airflow.models.param import Param +from orchestrate.utils import datacoves_utils + +@dag( + doc_md = __doc__, + catchup = False, + + default_args = datacoves_utils.set_default_args( + owner = "Ian", + owner_email = "ian@example.com" + ), + + description="Sample DAG showing templated overrides on the dbt decorator", + schedule = datacoves_utils.set_schedule(None), + + params = { + "dbt_command": Param( + "dbt build --select models/l1/schema_name+", + type="string", + description="dbt command to run", + ), + "schema_override": Param( + "MANUAL_REFRESH", + type="string", + description="Schema to build into for this run", + ), + }, + + tags=["transform", "parameters"], +) +def dbt_templated_overrides(): + + @task.datacoves_dbt( + connection_id="main_key_pair", + overrides={"schema": "{{ params.schema_override }}"}, + ) + def run_dbt(dbt_command): + return dbt_command + + run_dbt("{{ params.dbt_command }}") + +dbt_templated_overrides() diff --git a/orchestrate/dags/other_examples/retry_dbt_failures.py b/orchestrate/dags/other_examples/retry_dbt_failures.py index 7aec0626d..29d6bf604 100644 --- a/orchestrate/dags/other_examples/retry_dbt_failures.py +++ b/orchestrate/dags/other_examples/retry_dbt_failures.py @@ -12,7 +12,7 @@ default_args=datacoves_utils.set_default_args( owner = "Noel Gomez", - owner_email = "noel@example.com" + owner_email = "gomezn@example.com" ), schedule = datacoves_utils.set_schedule("0 0 1 */12 *"), diff --git a/orchestrate/dags/other_examples/test_oom.py b/orchestrate/dags/other_examples/test_oom.py new file mode 100644 index 000000000..65b683983 --- /dev/null +++ b/orchestrate/dags/other_examples/test_oom.py @@ -0,0 +1,62 @@ +""" +## OOM Test DAG +Deliberately allocates ~5GB of memory in a single task, then holds it for +2 minutes before releasing. Used to test that OOM detection/alerting pops up a +message when a task/pod runs out of memory. + +Trigger manually - not scheduled. + +Note: whether this actually gets OOM-killed depends on the memory +limit configured for the worker pod. If the pod's limit is below the +`target_gb` below, the task will be killed while allocating. If the limit +is higher, the task will succeed after holding the memory for `hold_seconds`. +""" +from pendulum import datetime + +from airflow.decorators import dag, task + +GB = 1024 ** 3 + +default_args = { + "start_date": datetime(2024, 1, 1), + "owner": "Fernando Mercado", + "email_on_failure": False, + "retries": 1, +} + + +@dag( + doc_md=__doc__, + catchup=False, + schedule=None, + default_args=default_args, + tags=["sample", "maintenance"], + dag_id="test_oom", +) +def test_oom(): + + @task + def consume_memory(target_gb: int = 5, hold_seconds: int = 120): + import time + + chunk_size = 256 * 1024 * 1024 # 256MB per chunk + target_bytes = target_gb * GB + + chunks = [] + allocated = 0 + + while allocated < target_bytes: + # bytearray() zero-fills on creation, which forces the pages + # to become resident rather than just reserved + chunks.append(bytearray(chunk_size)) + allocated += chunk_size + print(f"Allocated {allocated / GB:.2f} GB") + + print(f"Holding {allocated / GB:.2f} GB for {hold_seconds}s") + time.sleep(hold_seconds) + print("Releasing memory") + + consume_memory() + + +test_oom() diff --git a/orchestrate/dags_yml_definitions/daily_loan_run.yml b/orchestrate/dags_yml_definitions/daily_loan_run.yml index d8ab5c050..69016f4fc 100644 --- a/orchestrate/dags_yml_definitions/daily_loan_run.yml +++ b/orchestrate/dags_yml_definitions/daily_loan_run.yml @@ -1,8 +1,9 @@ doc_md: | - """ ## Sample DAG showing end-to-end ELT This DAG shows how to load data with 3 tools, then run dbt, then other tasks - """ + +imports: + - "from orchestrate.utils import datacoves_utils" description: "Loan Run" schedule: "0 0 1 */12 *" tags: @@ -35,15 +36,17 @@ nodes: tooltip: "dlt Extract and Load" tasks: - load_us_population: + load_loans_data: task_decorator: datacoves_bash - bash_command: "./load/dlt/load_data.py" + bash_command: "cd load/dlt && ./loans_data.py" + env: !py 'datacoves_utils.set_dlt_env_vars({"destinations": ["main_load_keypair"]})' + append_env: true transform: type: task task_decorator: datacoves_dbt - connection_id: main - bash_command: "dbt build -s 'tag:daily_run_airbyte+ tag:daily_run_fivetran+ -t prd'" + connection_id: main_key_pair + bash_command: "dbt build -s 'tag:daily_run_airbyte+ tag:daily_run_fivetran+'" dependencies: [ "extract_and_load_airbyte", @@ -54,6 +57,7 @@ nodes: marketing_automation: task_decorator: datacoves_bash type: task + bash_command: "echo 'send data to marketing tool'" dependencies: ["transform"] diff --git a/orchestrate/dags_yml_definitions/dbt_dag.yml b/orchestrate/dags_yml_definitions/dbt_dag.yml index 7502c39dc..7d22da85f 100644 --- a/orchestrate/dags_yml_definitions/dbt_dag.yml +++ b/orchestrate/dags_yml_definitions/dbt_dag.yml @@ -25,4 +25,10 @@ nodes: type: task task_decorator: datacoves_dbt connection_id: main_key_pair + # Task-level override: fail fast instead of inheriting the retries + # from default_args. + retries: 0 + # execution_timeout wants a timedelta, and !py is only resolved for + # DAG-level keys -- at node level dbt-coves emits it as a quoted + # string, which Airflow rejects. Left out until that is fixed. bash_command: "dbt debug" diff --git a/orchestrate/dags_yml_definitions/notifications_examples/yaml_slack_dag.yml b/orchestrate/dags_yml_definitions/notifications_examples/yaml_slack_dag.yml index 6cb3e06c8..167529610 100644 --- a/orchestrate/dags_yml_definitions/notifications_examples/yaml_slack_dag.yml +++ b/orchestrate/dags_yml_definitions/notifications_examples/yaml_slack_dag.yml @@ -29,7 +29,7 @@ notifications: nodes: transform: task_decorator: datacoves_dbt - connection_id: main + connection_id: main_key_pair type: task bash_command: "dbt run -s personal_loans" diff --git a/orchestrate/dags_yml_definitions/notifications_examples/yaml_teams_dag.yml b/orchestrate/dags_yml_definitions/notifications_examples/yaml_teams_dag.yml index d655717f8..3ab684c02 100644 --- a/orchestrate/dags_yml_definitions/notifications_examples/yaml_teams_dag.yml +++ b/orchestrate/dags_yml_definitions/notifications_examples/yaml_teams_dag.yml @@ -31,7 +31,7 @@ notifications: nodes: transform: task_decorator: datacoves_dbt - connection_id: main + connection_id: main_key_pair type: task bash_command: "dbt run -s personal_loans" diff --git a/secure/snowcap/apply.sh b/secure/snowcap/apply.sh index 0a6ddd8be..1b37e4530 100755 --- a/secure/snowcap/apply.sh +++ b/secure/snowcap/apply.sh @@ -75,7 +75,7 @@ if $USE_PII; then CONFIG_PATHS="--config $COMBINED_CONFIG" else ACCOUNT_TO_USE="$SNOWFLAKE_ACCOUNT" - EXCLUDE_RESOURCES="--exclude masking_policy,tag,tag_reference,tag_masking_policy_reference,row_access_policy" + EXCLUDE_RESOURCES="--exclude masking_policy,tag,tag_reference,tag_masking_policy_reference,row_access_policy,stream" USE_ACCOUNT_USAGE="" # Standard account - only include base resources CONFIG_PATHS="--config resources/" diff --git a/secure/snowcap/plan.sh b/secure/snowcap/plan.sh index 31ca6e95c..2dbafcacf 100755 --- a/secure/snowcap/plan.sh +++ b/secure/snowcap/plan.sh @@ -70,7 +70,7 @@ if $USE_PII; then CONFIG_PATHS="--config $COMBINED_CONFIG" else ACCOUNT_TO_USE="$SNOWFLAKE_ACCOUNT" - EXCLUDE_RESOURCES="--exclude masking_policy,tag,tag_reference,tag_masking_policy_reference,row_access_policy" + EXCLUDE_RESOURCES="--exclude masking_policy,tag,tag_reference,tag_masking_policy_reference,row_access_policy,stream" USE_ACCOUNT_USAGE="" # Standard account - only include base resources CONFIG_PATHS="--config resources/" diff --git a/secure/snowcap/resources/account.yml b/secure/snowcap/resources/account.yml index bd4ba5d41..9b612261e 100644 --- a/secure/snowcap/resources/account.yml +++ b/secure/snowcap/resources/account.yml @@ -10,3 +10,156 @@ account_parameters: # Still a Snowflake preview feature, gated behind account-wide preview access. - name: FEATURE_RBAC_INHERITED_GRANTS value: ENABLED + + + # Caps any single query at 1 hour so a runaway statement cannot burn credits + # Snowflake's default is 172800 seconds (2 days). + # + # Warehouses that set their own value in warehouses.yml, which override this. + - name: STATEMENT_TIMEOUT_IN_SECONDS + value: 600 + + # Snowflake's default is 0, meaning a statement waits behind a backlogged + # warehouse indefinitely. A silent stall rather than a failure. 10 minutes + # turns that into a visible error the orchestrator can retry. + - name: STATEMENT_QUEUED_TIMEOUT_IN_SECONDS + value: 600 + + # Snowflake's default is FALSE: when a client disconnects, its query keeps + # running and keeps the warehouse awake until STATEMENT_TIMEOUT fires. + # Already TRUE in the account but undeclared; this pins it. + - name: ABORT_DETACHED_QUERY + value: true + + # Snowflake's default is 10.0 TB. A guard against one ALTER DATABASE ... REFRESH + # replicating far more than intended. Declared at the default so the ceiling is + # a recorded decision rather than an inherited one. + - name: INITIAL_REPLICATION_SIZE_LIMIT_IN_TB + value: 10.0 + + # Snowflake's default is 0 -- no floor, so any database, schema or table can + # set DATA_RETENTION_TIME_IN_DAYS = 0 and disable Time Travel entirely. + # Effective retention becomes MAX(object value, this), so 1 guarantees a day + # of UNDROP on everything. + # + # 1 is also the ceiling on Standard Edition. On the Enterprise PII account + # this could go as high as 90; raising it there means moving the parameter out + # of this shared file first. + - name: MIN_DATA_RETENTION_TIME_IN_DAYS + value: 1 + + # Snowflake's default is FALSE, which lets anyone who can read a table run + # COPY INTO against a bucket they control with credentials written inline. + # TRUE forces unloads through a named stage. + - name: PREVENT_UNLOAD_TO_INLINE_URL + value: true + + # Snowflake's default is FALSE, allowing new external stages to carry cloud + # keys in the CREATE STAGE statement where they land in query history. + # Only affects stages created from here on, so existing stages keep working. + - name: REQUIRE_STORAGE_INTEGRATION_FOR_STAGE_CREATION + value: true + + # The matching operation-time rule stays OFF. This was tested against the + # live account on 2026-08-21, not assumed: + # + # ALTER ACCOUNT SET REQUIRE_STORAGE_INTEGRATION_FOR_STAGE_OPERATION = TRUE; + # LIST @RAW.RAW.EXT_JSONFILES_STAGE; + # -> 003160 (42601): Usage of stages with direct credentials, including + # public storage locations, has been forbidden. + # + # Note "including public storage locations". The parameter is not limited to + # stages carrying inline credentials -- it also blocks anonymous reads of + # public buckets. RAW.RAW.EXT_JSONFILES_STAGE has no STAGE_CREDENTIALS and + # points at the public s3://pslsnow1/jsonFiles/, and it is still blocked. + # + # That stage is live: 9 files, 126 queries in the last year, most recently + # 2026-08-19. Enabling this without migrating it first breaks that load. + # + # To enable, do one of these first: + # a) Give it a storage integration -- an IAM role in our AWS account with + # s3:GetObject/s3:ListBucket on the bucket, then + # ALTER STAGE ... SET STORAGE_INTEGRATION, declared here as + # s3_storage_integrations (snowcap supports the resource). + # b) Copy the 9 files into an internal stage and repoint the consumer. + # Simpler, given pslsnow1 is a third-party bucket we do not own. + # + # GREAT_BAY_DEV.PUBLIC.RAW -> s3://datacoves-sample-data-public/ was the + # other blocker. Zero references in 12 months of QUERY_HISTORY, no pipes, no + # external tables. Dropped 2026-08-21. + # + # - name: REQUIRE_STORAGE_INTEGRATION_FOR_STAGE_OPERATION + # value: true + + # Snowflake's default is FALSE, which lets any user grant privileges straight + # to another user, invisible to the role graph and to this repo. That defeats + # the z_* atomic-role model. Existing user grants are unaffected. + - name: DISABLE_USER_PRIVILEGE_GRANTS + value: true + + # Snowflake's default is FALSE. Already TRUE in the account it stops an MFA + # prompt on every new connection. Declared so the choice is recorded. + - name: ALLOW_CLIENT_MFA_CACHING + value: true + + # PERIODIC_DATA_REKEYING (yearly re-encryption with fresh keys) is worth + # enabling, but it is Enterprise Edition or higher and this file also applies + # to the standard account. Set it on the PII account separately. + + # Snowflake's default is America/Los_Angeles, so CURRENT_DATE rolls over on + # Pacific time. UTC is the only timezone nobody has to reason about. + # + # This changes what existing queries return: anything comparing timestamps + # against CURRENT_DATE, or using CONVERT_TIMEZONE without an explicit source, + # shifts by up to 8 hours on the first run after this applies. + - name: TIMEZONE + value: UTC + + # Snowflake's default is FALSE: an UPDATE whose join matches multiple source + # rows picks one and completes silently. MERGE already defaults to raising in + # the same situation this aligns the two. + - name: ERROR_ON_NONDETERMINISTIC_UPDATE + value: true + + # Declared at Snowflake's own default. What a bare TIMESTAMP column resolves + # to is worth pinning, because changing it later reinterprets existing DDL. + - name: TIMESTAMP_TYPE_MAPPING + value: TIMESTAMP_NTZ + + # Both declared at Snowflake's defaults: 0 means legacy ISO-like week + # semantics rather than a week anyone chose. Pinned so week-based reporting + # cannot start disagreeing after an account-level change nobody made here. + - name: WEEK_START + value: 0 + - name: WEEK_OF_YEAR_POLICY + value: 0 + + # Snowflake's default is 10 consecutive failures before a task suspends + # itself. If each attempt resumes a warehouse, that is ten billable failures + # on a task that was never going to succeed. + - name: SUSPEND_TASK_AFTER_NUM_FAILURES + value: 3 + + # Declared at Snowflake's default of 1 hour, matching STATEMENT_TIMEOUT_IN_SECONDS + # above so the two stay deliberately in step. + - name: USER_TASK_TIMEOUT_MS + value: 3600000 + + # Both default to OFF, so nothing a Python UDF or stored procedure logs ever + # reaches the event table -- which you discover the first time you need it. + # WARN keeps the volume (and the storage cost) low; ON_EVENT records spans + # only where code explicitly emits them. + - name: LOG_LEVEL + value: WARN + - name: TRACE_LEVEL + value: ON_EVENT + + # MAX_CONCURRENCY_LEVEL is intentionally left at Snowflake's default of 8. + # It is a per-warehouse tuning knob, and lowering it without evidence of + # contention costs throughput. Set it in warehouses.yml on a specific + # warehouse once queuing shows up in QUERY_HISTORY. + + # NETWORK_POLICY is intentionally not set here. It needs a real IP allowlist, + # and there is no safe default -- a wrong one locks every user and service + # account out of the account. Create a network_policy resource with the + # actual office/VPN/service CIDRs, then attach it. diff --git a/secure/snowcap/resources/object_templates/schema.yml b/secure/snowcap/resources/object_templates/schema.yml index 30023f34b..edb9f4ccc 100644 --- a/secure/snowcap/resources/object_templates/schema.yml +++ b/secure/snowcap/resources/object_templates/schema.yml @@ -6,6 +6,14 @@ schemas: owner: "{{ each.value.get('owner', parent.owner) }}" managed_access: true + # BALBOA_QA companion schemas share the BALBOA schema role below. + - for_each: var.schemas + where: "each.value.name.split('.')[0] == 'BALBOA'" + name: "{{ each.value.name.split('.')[1] }}" + database: BALBOA_QA + owner: "{{ each.value.get('owner', parent.owner) }}" + managed_access: true + # Schema roles roles: - for_each: var.schemas @@ -17,3 +25,10 @@ grants: priv: USAGE on: "schema {{ each.value.name }}" to: "z_schema__{{ each.value.name.split('.')[1] }}" + + # QA schema grants + - for_each: var.schemas + where: "each.value.name.split('.')[0] == 'BALBOA'" + priv: USAGE + on: "schema BALBOA_QA.{{ each.value.name.split('.')[1] }}" + to: "z_schema__{{ each.value.name.split('.')[1] }}" diff --git a/secure/snowcap/resources/roles__functional.yml b/secure/snowcap/resources/roles__functional.yml index d1bb48588..46986d512 100644 --- a/secure/snowcap/resources/roles__functional.yml +++ b/secure/snowcap/resources/roles__functional.yml @@ -44,6 +44,9 @@ role_grants: - to_role: finance_team roles: - z_db__balboa + - z_db__balboa_qa + - z_tables_views__select + - z_schema__l3_accounts_payable - z_schema__l3_loan_analytics - z_wh__wh_finance @@ -63,13 +66,18 @@ role_grants: - z_schema__l1_google_analytics_4 - z_schema__l1_loans - z_schema__l1_observe + - z_schema__l1_erp - z_schema__l1_usgs__earthquake_data - z_schema__l1_us_population - z_schema__l2_country_demographics - z_schema__l2_covid_observations + - z_schema__l2_invoices + - z_schema__l2_purchase_orders - z_schema__l2_snowflake_usage + - z_schema__l2_vendors + - z_schema__l3_accounts_payable - z_schema__l3_covid_analytics - z_schema__l3_earthquake_analytics - z_schema__l3_loan_analytics diff --git a/secure/snowcap/resources/schemas.yml b/secure/snowcap/resources/schemas.yml index e5ad0ff45..e24c509f3 100644 --- a/secure/snowcap/resources/schemas.yml +++ b/secure/snowcap/resources/schemas.yml @@ -23,13 +23,18 @@ vars: - name: BALBOA.L1_GOOGLE_ANALYTICS_4 - name: BALBOA.L1_LOANS - name: BALBOA.L1_OBSERVE + - name: BALBOA.L1_ERP - name: BALBOA.L1_USGS__EARTHQUAKE_DATA - name: BALBOA.L1_US_POPULATION - name: BALBOA.L2_COUNTRY_DEMOGRAPHICS - name: BALBOA.L2_COVID_OBSERVATIONS + - name: BALBOA.L2_INVOICES + - name: BALBOA.L2_PURCHASE_ORDERS - name: BALBOA.L2_SNOWFLAKE_USAGE + - name: BALBOA.L2_VENDORS + - name: BALBOA.L3_ACCOUNTS_PAYABLE - name: BALBOA.L3_COVID_ANALYTICS - name: BALBOA.L3_EARTHQUAKE_ANALYTICS - name: BALBOA.L3_LOAN_ANALYTICS @@ -42,23 +47,3 @@ vars: # GREAT_BAY DB - name: GREAT_BAY.COVE_MARKETING -grants: - # BALBOA_QA is a clone of BALBOA (transform/macros/tooling/blue-green). - # CREATE DATABASE ... CLONE copies grants on child objects, so every - # z_schema__ role already reaches the QA copy of its schema -- these - # grants exist in Snowflake whether or not we declare them. Declaring them - # keeps sync from revoking them. - # - # Deliberately per-schema rather than "all schemas in database balboa_qa": - # roles are scoped by layer (finance_team sees only L3_LOAN_ANALYTICS), and a - # database-wide grant would hand every such role the whole QA database. - # - # The QA schemas themselves stay undeclared -- the clone creates them, and - # the `where` keeps this off schemas in RAW, GOVERNANCE and the rest, which - # have no QA counterpart. - - for_each: var.schemas - where: "each.value.name.split('.')[0] == 'BALBOA'" - priv: "USAGE" - on: "schema BALBOA_QA.{{ each.value.name.split('.')[1] }}" - to: "z_schema__{{ each.value.name.split('.')[1] }}" - diff --git a/secure/snowcap/resources/warehouses.yml b/secure/snowcap/resources/warehouses.yml index be8d50ebe..3fe3c1a1b 100644 --- a/secure/snowcap/resources/warehouses.yml +++ b/secure/snowcap/resources/warehouses.yml @@ -5,7 +5,7 @@ vars: - name: wh_admin size: x-small auto_suspend: 60 # seconds 0 means it never suspends - statement_timeout_in_seconds: 7200 # seconds, defaults to 3600 if not defined + statement_timeout_in_seconds: 3600 # seconds; object_templates/warehouses.yml applies 3600 when omitted - name: wh_catalog size: x-small auto_suspend: 60 diff --git a/transform/dbt_project.yml b/transform/dbt_project.yml index 7f1b2044a..2394bc732 100644 --- a/transform/dbt_project.yml +++ b/transform/dbt_project.yml @@ -74,6 +74,8 @@ models: +schema: L1_COVID19_EPIDEMIOLOGICAL_DATA loans: +schema: L1_LOANS + erp: + +schema: L1_ERP observe: +schema: L1_OBSERVE us_population: @@ -91,6 +93,12 @@ models: +schema: L2_COVID_OBSERVATIONS snowflake_usage: +schema: L2_SNOWFLAKE_USAGE + vendors: + +schema: L2_VENDORS + invoices: + +schema: L2_INVOICES + purchase_orders: + +schema: L2_PURCHASE_ORDERS L3_coves: +group: marketing @@ -103,6 +111,8 @@ models: +schema: L3_EARTHQUAKE_ANALYTICS loan_analytics: +schema: L3_LOAN_ANALYTICS + accounts_payable: + +schema: L3_ACCOUNTS_PAYABLE # cannot persist docs on dynamic tables # +persist_docs: # relation: false @@ -153,6 +163,8 @@ data_tests: vars: 'dbt_date:time_zone': 'America/Los_Angeles' + authored_history_through: '2025-12-01' + as_of: '2026-09-08' # Snowcap governance macros snowcap_tag_database: "GOVERNANCE" diff --git a/transform/macros/spine.sql b/transform/macros/spine.sql new file mode 100644 index 000000000..85ea8ec4e --- /dev/null +++ b/transform/macros/spine.sql @@ -0,0 +1,15 @@ +{% macro trailing_edge_month() %} + {% set as_of = var('as_of') %} + date_trunc('month', dateadd(month, -1, to_date('{{ as_of }}'))) +{% endmacro %} + +{% macro procurement_month_spine() %} + {% set trailing_edge = trailing_edge_month() %} + + select + cast( + dateadd(month, seq4(), dateadd(month, -23, {{ trailing_edge }})) + as date + ) as month_start + from table(generator(rowcount => 24)) +{% endmacro %} diff --git a/transform/models/L1_inlets/erp/stg_invoices.sql b/transform/models/L1_inlets/erp/stg_invoices.sql new file mode 100644 index 000000000..4b333267f --- /dev/null +++ b/transform/models/L1_inlets/erp/stg_invoices.sql @@ -0,0 +1,70 @@ +with purchase_orders as ( + + select + po_id, + vendor_id, + po_date, + po_amount + from {{ ref('stg_purchase_orders') }} + +), + +vendors as ( + + select + vendor_id, + vendor_seq, + activity_bucket, + invoice_po_through_month, + invoice_lag_months + from {{ ref('stg_vendors') }} + +), + +candidate_invoices as ( + + select + purchase_orders.po_id, + purchase_orders.vendor_id, + purchase_orders.po_amount, + vendors.vendor_seq, + cast(dateadd(month, vendors.invoice_lag_months, purchase_orders.po_date) as date) as invoice_date + from purchase_orders + inner join vendors + on purchase_orders.vendor_id = vendors.vendor_id + where purchase_orders.po_date <= vendors.invoice_po_through_month + and ( + vendors.activity_bucket <> 'C' + or purchase_orders.po_date = vendors.invoice_po_through_month + ) + +), + +invoices_in_span as ( + + select + po_id, + vendor_id, + po_amount, + vendor_seq, + invoice_date + from candidate_invoices + where invoice_date <= {{ trailing_edge_month() }} + +) + +select + replace(po_id, '-PO-', '-INV-')::varchar as invoice_id, + po_id, + vendor_id, + invoice_date, + po_amount::bigint as invoice_amount, + case + when mod(vendor_seq + month(invoice_date), 5) = 0 then null + else dateadd(day, 15, invoice_date)::date + end as paid_date, + case + when mod(vendor_seq + month(invoice_date), 5) = 0 then 'open' + else 'paid' + end::varchar as invoice_status +from invoices_in_span diff --git a/transform/models/L1_inlets/erp/stg_invoices.yml b/transform/models/L1_inlets/erp/stg_invoices.yml new file mode 100644 index 000000000..d9a420eb9 --- /dev/null +++ b/transform/models/L1_inlets/erp/stg_invoices.yml @@ -0,0 +1,46 @@ +version: 2 + +models: + - name: stg_invoices + description: Deterministic ERP invoices, including delayed and missing invoice scenarios. + columns: + - name: invoice_id + description: Invoice primary key. + data_tests: [not_null, unique] + - name: po_id + description: Billed purchase order foreign key. + data_tests: + - not_null + - unique + - relationships: + arguments: + to: ref('stg_purchase_orders') + field: po_id + - name: vendor_id + description: Vendor foreign key. + data_tests: + - not_null + - relationships: + arguments: + to: ref('stg_vendors') + field: vendor_id + - name: invoice_date + description: Invoice activity date. + data_tests: [not_null] + - name: invoice_amount + description: Positive whole-currency invoice amount equal to the PO amount. + data_tests: + - not_null + - dbt_utils.expression_is_true: + arguments: + expression: '> 0' + - name: paid_date + description: Payment date; null while the invoice remains open. + + - name: invoice_status + description: Payment state for the invoice. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['paid', 'open'] diff --git a/transform/models/L1_inlets/erp/stg_purchase_orders.sql b/transform/models/L1_inlets/erp/stg_purchase_orders.sql new file mode 100644 index 000000000..843ec749d --- /dev/null +++ b/transform/models/L1_inlets/erp/stg_purchase_orders.sql @@ -0,0 +1,41 @@ +with month_spine as ( + + {{ procurement_month_spine() }} + +), + +vendors as ( + + select + vendor_id, + category, + onboarded_date, + po_active_through_month, + vendor_seq + from {{ ref('stg_vendors') }} + +) + +select + concat(vendors.vendor_id, '-PO-', to_char(month_spine.month_start, 'YYYYMM'))::varchar as po_id, + vendors.vendor_id, + month_spine.month_start::date as po_date, + vendors.category, + ( + case vendors.category + when 'CRO Services' then 24000 + when 'Clinical Reagents' then 18000 + when 'IT / Software' then 12000 + when 'Logistics' then 9000 + when 'Lab Supplies' then 15000 + when 'Packaging' then 7000 + when 'Consulting' then 16000 + when 'Facilities' then 11000 + end + + mod(vendors.vendor_seq, 4) * 125 + )::bigint as po_amount +from vendors +inner join month_spine + on month_spine.month_start >= vendors.onboarded_date + and month_spine.month_start <= vendors.po_active_through_month +where mod(datediff(month, vendors.onboarded_date, month_spine.month_start), 3) = 0 diff --git a/transform/models/L1_inlets/erp/stg_purchase_orders.yml b/transform/models/L1_inlets/erp/stg_purchase_orders.yml new file mode 100644 index 000000000..8062428e6 --- /dev/null +++ b/transform/models/L1_inlets/erp/stg_purchase_orders.yml @@ -0,0 +1,30 @@ +version: 2 + +models: + - name: stg_purchase_orders + description: Deterministic quarterly ERP purchase orders at monthly date grain. + columns: + - name: po_id + description: Purchase order primary key. + data_tests: [not_null, unique] + - name: vendor_id + description: Vendor foreign key. + data_tests: + - not_null + - relationships: + arguments: + to: ref('stg_vendors') + field: vendor_id + - name: po_date + description: First day of the purchase order month. + data_tests: [not_null] + - name: category + description: Vendor category carried onto the purchase order. + data_tests: [not_null] + - name: po_amount + description: Positive committed spend in whole currency units. + data_tests: + - not_null + - dbt_utils.expression_is_true: + arguments: + expression: '> 0' diff --git a/transform/models/L1_inlets/erp/stg_vendors.sql b/transform/models/L1_inlets/erp/stg_vendors.sql new file mode 100644 index 000000000..c45bfb9f5 --- /dev/null +++ b/transform/models/L1_inlets/erp/stg_vendors.sql @@ -0,0 +1,106 @@ +with month_spine as ( + + {{ procurement_month_spine() }} + +), + +authored_vendors as ( + + select + vendor_id::varchar as vendor_id, + vendor_name::varchar as vendor_name, + category::varchar as category, + status::varchar as status, + onboarded_date::date as onboarded_date + from {{ ref('seed_vendors') }} + +), + +archetypes as ( + + select + archetype_index::integer as archetype_index, + vendor_name::varchar as vendor_name, + category::varchar as category, + status::varchar as status + from {{ ref('seed_vendor_archetypes') }} + +), + +generated_months as ( + + select + month_start, + row_number() over (order by month_start) as generated_vendor_number + from month_spine + where month_start > to_date('{{ var("authored_history_through") }}') + +), + +generated_vendors as ( + + select + concat('VENT', to_char(generated_months.month_start, 'YYYYMM'))::varchar as vendor_id, + concat(archetypes.vendor_name, ' ', to_char(generated_months.month_start, 'YYYYMM'))::varchar as vendor_name, + archetypes.category, + archetypes.status, + generated_months.month_start::date as onboarded_date + from generated_months + join archetypes + on archetypes.archetype_index = mod( + generated_months.generated_vendor_number - 1, + (select count(*) from archetypes) + ) + +), + +all_vendors as ( + + select * from authored_vendors + union all + select * from generated_vendors + +), + +bucketed as ( + + select + vendor_id, + vendor_name, + category, + status, + onboarded_date, + row_number() over (order by vendor_id) as vendor_seq, + case + when vendor_id like 'VENT%' then 'A' + when to_number(replace(vendor_id, 'VEN', '')) <= 6 then 'A' + when to_number(replace(vendor_id, 'VEN', '')) <= 11 then 'B' + when to_number(replace(vendor_id, 'VEN', '')) <= 17 then 'C' + else 'D' + end::varchar as activity_bucket + from all_vendors + +) + +select + vendor_id, + vendor_name, + category, + status, + onboarded_date, + vendor_seq, + activity_bucket, + case + when activity_bucket in ('A', 'B') then {{ trailing_edge_month() }} + else dateadd(month, -15, {{ trailing_edge_month() }}) + end::date as po_active_through_month, + case + when activity_bucket = 'B' then dateadd(month, -15, {{ trailing_edge_month() }}) + when activity_bucket in ('A', 'B') then {{ trailing_edge_month() }} + else dateadd(month, -15, {{ trailing_edge_month() }}) + end::date as invoice_po_through_month, + case + when activity_bucket = 'C' then 14 + else 1 + end::integer as invoice_lag_months +from bucketed diff --git a/transform/models/L1_inlets/erp/stg_vendors.yml b/transform/models/L1_inlets/erp/stg_vendors.yml new file mode 100644 index 000000000..b54fd4493 --- /dev/null +++ b/transform/models/L1_inlets/erp/stg_vendors.yml @@ -0,0 +1,44 @@ +version: 2 + +models: + - name: stg_vendors + description: ERP vendors with deterministic activity-generation drivers. + columns: + - name: vendor_id + description: Vendor primary key. + data_tests: [not_null, unique] + - name: vendor_name + description: Vendor display name. + data_tests: [not_null] + - name: category + description: Vendor spend category. + data_tests: [not_null] + - name: status + description: Recorded lifecycle flag, independent from behavioral activity. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['active', 'inactive'] + - name: onboarded_date + description: First day of the vendor onboarding month. + data_tests: [not_null] + - name: activity_bucket + description: Activity profile controlling PO and invoice recency. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['A', 'B', 'C', 'D'] + - name: vendor_seq + description: Stable sequence used to assign deterministic activity profiles. + data_tests: [not_null] + - name: po_active_through_month + description: Final month in which the vendor can generate purchase orders. + data_tests: [not_null] + - name: invoice_po_through_month + description: Final purchase order month eligible to generate an invoice. + data_tests: [not_null] + - name: invoice_lag_months + description: Whole-month lag between an eligible purchase order and its invoice. + data_tests: [not_null] diff --git a/transform/models/L2_bays/invoices/fct_invoices.sql b/transform/models/L2_bays/invoices/fct_invoices.sql new file mode 100644 index 000000000..ff29ae7e3 --- /dev/null +++ b/transform/models/L2_bays/invoices/fct_invoices.sql @@ -0,0 +1,9 @@ +select + invoice_id, + po_id, + vendor_id, + invoice_date, + invoice_amount, + paid_date, + invoice_status +from {{ ref('stg_invoices') }} diff --git a/transform/models/L2_bays/invoices/fct_invoices.yml b/transform/models/L2_bays/invoices/fct_invoices.yml new file mode 100644 index 000000000..715975ad1 --- /dev/null +++ b/transform/models/L2_bays/invoices/fct_invoices.yml @@ -0,0 +1,45 @@ +version: 2 + +models: + - name: fct_invoices + description: Reusable ERP invoice entity at one row per invoice. + columns: + - name: invoice_id + description: Invoice primary key. + data_tests: [not_null, unique] + - name: po_id + description: Billed purchase order foreign key. + data_tests: + - not_null + - unique + - relationships: + arguments: + to: ref('fct_purchase_orders') + field: po_id + - name: vendor_id + description: Vendor foreign key. + data_tests: + - not_null + - relationships: + arguments: + to: ref('dim_vendors') + field: vendor_id + - name: invoice_date + description: Invoice activity date. + data_tests: [not_null] + - name: invoice_amount + description: Positive whole-currency invoice amount equal to the PO amount. + data_tests: + - not_null + - dbt_utils.expression_is_true: + arguments: + expression: '> 0' + - name: paid_date + description: Null while the invoice is open. + - name: invoice_status + description: Payment state for the invoice. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['paid', 'open'] diff --git a/transform/models/L2_bays/purchase_orders/fct_purchase_orders.sql b/transform/models/L2_bays/purchase_orders/fct_purchase_orders.sql new file mode 100644 index 000000000..f88f735ab --- /dev/null +++ b/transform/models/L2_bays/purchase_orders/fct_purchase_orders.sql @@ -0,0 +1,7 @@ +select + po_id, + vendor_id, + po_date, + category, + po_amount +from {{ ref('stg_purchase_orders') }} diff --git a/transform/models/L2_bays/purchase_orders/fct_purchase_orders.yml b/transform/models/L2_bays/purchase_orders/fct_purchase_orders.yml new file mode 100644 index 000000000..3593fb1ef --- /dev/null +++ b/transform/models/L2_bays/purchase_orders/fct_purchase_orders.yml @@ -0,0 +1,30 @@ +version: 2 + +models: + - name: fct_purchase_orders + description: Reusable ERP purchase order entity at one row per purchase order. + columns: + - name: po_id + description: Purchase order primary key. + data_tests: [not_null, unique] + - name: vendor_id + description: Vendor foreign key. + data_tests: + - not_null + - relationships: + arguments: + to: ref('dim_vendors') + field: vendor_id + - name: po_date + description: First day of the purchase order month. + data_tests: [not_null] + - name: category + description: Vendor category carried onto the purchase order. + data_tests: [not_null] + - name: po_amount + description: Positive committed spend in whole currency units. + data_tests: + - not_null + - dbt_utils.expression_is_true: + arguments: + expression: '> 0' diff --git a/transform/models/L2_bays/vendors/dim_vendors.sql b/transform/models/L2_bays/vendors/dim_vendors.sql new file mode 100644 index 000000000..e7470dc9c --- /dev/null +++ b/transform/models/L2_bays/vendors/dim_vendors.sql @@ -0,0 +1,7 @@ +select + vendor_id, + vendor_name, + category, + status, + onboarded_date +from {{ ref('stg_vendors') }} diff --git a/transform/models/L2_bays/vendors/dim_vendors.yml b/transform/models/L2_bays/vendors/dim_vendors.yml new file mode 100644 index 000000000..27403bb4c --- /dev/null +++ b/transform/models/L2_bays/vendors/dim_vendors.yml @@ -0,0 +1,25 @@ +version: 2 + +models: + - name: dim_vendors + description: Reusable ERP vendor entity at one row per vendor. + columns: + - name: vendor_id + description: Vendor primary key and fact join target. + data_tests: [not_null, unique] + - name: vendor_name + description: Vendor display name. + data_tests: [not_null] + - name: category + description: Vendor spend category. + data_tests: [not_null] + - name: status + description: Recorded lifecycle flag, not a behavioral activity measure. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['active', 'inactive'] + - name: onboarded_date + description: First day of the vendor onboarding month. + data_tests: [not_null] diff --git a/transform/models/L3_coves/accounts_payable/mart_active_vendors.sql b/transform/models/L3_coves/accounts_payable/mart_active_vendors.sql new file mode 100644 index 000000000..eed007312 --- /dev/null +++ b/transform/models/L3_coves/accounts_payable/mart_active_vendors.sql @@ -0,0 +1,70 @@ +{{ config(materialized='table') }} + +with purchase_orders as ( + + select + vendor_id, + po_date as activity_date, + 'purchase_order' as activity_source + from {{ ref('fct_purchase_orders') }} + +), + +invoices as ( + + select + vendor_id, + invoice_date as activity_date, + 'invoice' as activity_source + from {{ ref('fct_invoices') }} + +), + +activity as ( + + select * from purchase_orders + union all + select * from invoices + +), + +bounds as ( + + select max(activity_date) as reference_date + from activity + +) + +select + vendors.vendor_id, + vendors.vendor_name, + vendors.category, + vendors.status as recorded_status, + bounds.reference_date, + max(case when activity.activity_source = 'purchase_order' then activity.activity_date end) + as last_po_date, + max(case when activity.activity_source = 'invoice' then activity.activity_date end) + as last_invoice_date, + count_if( + activity.activity_source = 'purchase_order' + and activity.activity_date > dateadd(month, -12, bounds.reference_date) + and activity.activity_date <= bounds.reference_date + ) as trailing_12_month_po_count, + count_if( + activity.activity_source = 'invoice' + and activity.activity_date > dateadd(month, -12, bounds.reference_date) + and activity.activity_date <= bounds.reference_date + ) as trailing_12_month_invoice_count, + ( + trailing_12_month_po_count + trailing_12_month_invoice_count > 0 + ) as is_active +from {{ ref('dim_vendors') }} as vendors +cross join bounds +left join activity + on vendors.vendor_id = activity.vendor_id +group by + vendors.vendor_id, + vendors.vendor_name, + vendors.category, + vendors.status, + bounds.reference_date diff --git a/transform/models/L3_coves/accounts_payable/mart_active_vendors.yml b/transform/models/L3_coves/accounts_payable/mart_active_vendors.yml new file mode 100644 index 000000000..c26b73e1b --- /dev/null +++ b/transform/models/L3_coves/accounts_payable/mart_active_vendors.yml @@ -0,0 +1,38 @@ +version: 2 + +models: + - name: mart_active_vendors + description: Accounts Payable cove for behavioral Active Vendor analysis using trailing-12-month PO-or-invoice activity. + columns: + - name: vendor_id + description: Vendor primary key. + data_tests: [not_null, unique] + - name: vendor_name + description: Vendor display name. + data_tests: [not_null] + - name: category + description: Vendor spend category. + data_tests: [not_null] + - name: recorded_status + description: Stored ERP lifecycle flag, independent from behavioral activity. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['active', 'inactive'] + - name: reference_date + description: Latest PO or invoice date in the fixture. + data_tests: [not_null] + - name: last_po_date + description: Latest purchase order date for the vendor. + - name: last_invoice_date + description: Latest invoice date for the vendor. + - name: trailing_12_month_po_count + description: Count of purchase orders in the behavioral activity window. + data_tests: [not_null] + - name: trailing_12_month_invoice_count + description: Count of invoices in the behavioral activity window. + data_tests: [not_null] + - name: is_active + description: True when the vendor has at least one PO or invoice in the trailing 12 months. + data_tests: [not_null] diff --git a/transform/seeds/seed_vendor_archetypes.csv b/transform/seeds/seed_vendor_archetypes.csv new file mode 100644 index 000000000..8356fd853 --- /dev/null +++ b/transform/seeds/seed_vendor_archetypes.csv @@ -0,0 +1,5 @@ +archetype_index,vendor_name,category,status +0,Northstar Bioanalytics,CRO Services,active +1,Armitage Cloud Systems,IT / Software,active +2,Bellwether Packaging Co,Packaging,inactive +3,Kestrel Facilities Group,Facilities,active diff --git a/transform/seeds/seed_vendors.csv b/transform/seeds/seed_vendors.csv new file mode 100644 index 000000000..a6afb27b3 --- /dev/null +++ b/transform/seeds/seed_vendors.csv @@ -0,0 +1,27 @@ +vendor_id,vendor_name,category,status,onboarded_date +VEN01,Acadia Clinical Research,CRO Services,inactive,2024-09-01 +VEN02,Bluewater Trial Logistics,Logistics,active,2024-09-01 +VEN03,Caldera Lab Supplies,Lab Supplies,active,2024-09-01 +VEN04,Delphi Data Systems,IT / Software,active,2024-09-01 +VEN05,Evergreen Consulting Partners,Consulting,active,2024-09-01 +VEN06,Fairhaven Packaging,Packaging,active,2024-09-01 +VEN07,Granite Clinical Reagents,Clinical Reagents,inactive,2024-09-01 +VEN08,Harbor Point Facilities,Facilities,active,2024-09-01 +VEN09,Ironwood Logistics Group,Logistics,active,2024-09-01 +VEN10,Juniper Software Labs,IT / Software,active,2024-09-01 +VEN11,Keystone Research Services,CRO Services,active,2024-09-01 +VEN12,Lakeshore Reagents Co,Clinical Reagents,inactive,2024-11-01 +VEN13,Meridian Lab Equipment,Lab Supplies,active,2024-11-01 +VEN14,Norwood Packaging Works,Packaging,active,2024-11-01 +VEN15,Oakridge Advisory Group,Consulting,active,2024-11-01 +VEN16,Parkside Facilities Mgmt,Facilities,active,2024-11-01 +VEN17,Quarry Clinical Research,CRO Services,active,2024-11-01 +VEN18,Riverside Logistics Network,Logistics,active,2024-09-01 +VEN19,Summit Software Solutions,IT / Software,inactive,2024-09-01 +VEN20,Tidewater Lab Supplies,Lab Supplies,inactive,2024-09-01 +VEN21,Union Packaging Partners,Packaging,inactive,2024-09-01 +VEN22,Valley Consulting Group,Consulting,inactive,2024-09-01 +VEN23,Westgate Facilities Services,Facilities,inactive,2024-09-01 +VEN24,Xenon Clinical Reagents,Clinical Reagents,inactive,2024-09-01 +VEN25,Yorktown Research Partners,CRO Services,inactive,2024-09-01 +VEN26,Zephyr Logistics Services,Logistics,inactive,2024-09-01 diff --git a/transform/seeds/seed_vendors.yml b/transform/seeds/seed_vendors.yml new file mode 100644 index 000000000..80ed3f3f1 --- /dev/null +++ b/transform/seeds/seed_vendors.yml @@ -0,0 +1,58 @@ +version: 2 + +seeds: + - name: seed_vendors + description: Curated procurement vendors used to demonstrate behavioral activity. + config: + column_types: + vendor_id: varchar + vendor_name: varchar + category: varchar + status: varchar + onboarded_date: date + columns: + - name: vendor_id + description: Vendor primary key. + data_tests: [not_null, unique] + - name: vendor_name + description: Vendor display name. + data_tests: [not_null] + - name: category + description: Vendor spend category. + data_tests: [not_null] + - name: status + description: Recorded lifecycle flag, independent from behavioral activity. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['active', 'inactive'] + - name: onboarded_date + description: First day of the vendor onboarding month. + data_tests: [not_null] + + - name: seed_vendor_archetypes + description: Dateless profiles used for generated monthly vendor onboardings. + config: + column_types: + archetype_index: integer + vendor_name: varchar + category: varchar + status: varchar + columns: + - name: archetype_index + description: Zero-based round-robin archetype position. + data_tests: [not_null, unique] + - name: vendor_name + description: Generated vendor base name. + data_tests: [not_null] + - name: category + description: Generated vendor spend category. + data_tests: [not_null] + - name: status + description: Recorded lifecycle flag for generated vendors. + data_tests: + - not_null + - accepted_values: + arguments: + values: ['active', 'inactive'] diff --git a/transform/tests/L3_coves/accounts_payable/active_vendor_reference_counts.sql b/transform/tests/L3_coves/accounts_payable/active_vendor_reference_counts.sql new file mode 100644 index 000000000..90e56cfad --- /dev/null +++ b/transform/tests/L3_coves/accounts_payable/active_vendor_reference_counts.sql @@ -0,0 +1,52 @@ +{{ config(tags=['requires_fixture_data']) }} + +with reference_dates as ( + + select to_date('2026-06-01') as reference_date, 32 as expected_total, 17 as expected_active + union all + select to_date('2026-07-01') as reference_date, 33 as expected_total, 24 as expected_active + union all + select to_date('2026-08-01') as reference_date, 34 as expected_total, 25 as expected_active + +), + +activity as ( + + select vendor_id, po_date as activity_date + from {{ ref('fct_purchase_orders') }} + + union all + + select vendor_id, invoice_date as activity_date + from {{ ref('fct_invoices') }} + +), + +actual_counts as ( + + select + reference_dates.reference_date, + count(distinct dim_vendors.vendor_id) as actual_total, + count(distinct activity.vendor_id) as actual_active + from reference_dates + left join {{ ref('dim_vendors') }} as dim_vendors + on dim_vendors.onboarded_date <= reference_dates.reference_date + left join activity + on dim_vendors.vendor_id = activity.vendor_id + and activity.activity_date > dateadd(month, -12, reference_dates.reference_date) + and activity.activity_date <= reference_dates.reference_date + group by reference_dates.reference_date + +) + +select + reference_dates.reference_date, + reference_dates.expected_total, + actual_counts.actual_total, + reference_dates.expected_active, + actual_counts.actual_active +from reference_dates +inner join actual_counts + on reference_dates.reference_date = actual_counts.reference_date +where reference_dates.expected_total <> actual_counts.actual_total + or reference_dates.expected_active <> actual_counts.actual_active