misc.python.materialize.zippy.copy_s3_actions

 1# Copyright Materialize, Inc. and contributors. All rights reserved.
 2#
 3# Use of this software is governed by the Business Source License
 4# included in the LICENSE file at the root of this repository.
 5#
 6# As of the Change Date specified in that file, in accordance with
 7# the Business Source License, use of this software will be governed
 8# by the Apache License, Version 2.0.
 9
10import random
11from textwrap import dedent
12
13from materialize.mzcompose.composition import Composition
14from materialize.zippy.balancerd_capabilities import BalancerdIsRunning
15from materialize.zippy.copy_s3_capabilities import S3ObjectExists
16from materialize.zippy.framework import Action, Capabilities, Capability, State
17from materialize.zippy.mz_capabilities import MzIsRunning
18from materialize.zippy.table_capabilities import TableExists
19
20
21class CopyToS3(Action):
22    """Performs a COPY TO S3 and records the resulting S3 object as a capability."""
23
24    @classmethod
25    def requires(cls) -> set[type[Capability]]:
26        return {BalancerdIsRunning, MzIsRunning, TableExists}
27
28    def __init__(self, capabilities: Capabilities) -> None:
29        self.table = random.choice(capabilities.get(TableExists))
30        super().__init__(capabilities)
31        self.s3_key = f"zippy/{self.seqno}"
32        self.s3_object = S3ObjectExists(
33            s3_key=self.s3_key,
34            min_val=self.table.watermarks.min,
35            max_val=self.table.watermarks.max,
36        )
37
38    def provides(self) -> list[Capability]:
39        return [self.s3_object]
40
41    def run(self, c: Composition, state: State) -> None:
42        conn_name = f"zippy_aws_conn_{self.seqno}"
43        secret_name = f"zippy_minio_secret_{self.seqno}"
44
45        c.testdrive(
46            dedent(f"""
47                > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}'
48                > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name})
49                > COPY (SELECT f1 FROM {self.table.name}) TO 's3://copytos3/{self.s3_key}' WITH (AWS CONNECTION = {conn_name}, FORMAT = 'csv')
50                > DROP CONNECTION {conn_name}
51                > DROP SECRET {secret_name}
52                """),
53            mz_service=state.mz_service,
54        )
55        self.s3_object.min_val = self.table.watermarks.min
56        self.s3_object.max_val = self.table.watermarks.max
57
58
59class CopyFromS3(Action):
60    """Reads a previously written S3 object back into a staging table and validates the row count."""
61
62    @classmethod
63    def requires(cls) -> set[type[Capability]]:
64        return {BalancerdIsRunning, MzIsRunning, S3ObjectExists}
65
66    def __init__(self, capabilities: Capabilities) -> None:
67        self.s3_object = random.choice(capabilities.get(S3ObjectExists))
68        super().__init__(capabilities)
69
70    def run(self, c: Composition, state: State) -> None:
71        conn_name = f"zippy_aws_conn_{self.seqno}"
72        secret_name = f"zippy_minio_secret_{self.seqno}"
73        staging_table = f"zippy_s3_staging_{self.seqno}"
74
75        c.testdrive(
76            dedent(f"""
77                > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}'
78                > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name})
79                > CREATE TABLE {staging_table} (f1 INTEGER)
80                > COPY INTO {staging_table} FROM 's3://copytos3/{self.s3_object.s3_key}' (FORMAT CSV, AWS CONNECTION = {conn_name})
81                > SELECT MIN(f1), MAX(f1), COUNT(f1), COUNT(DISTINCT f1) FROM {staging_table}
82                {self.s3_object.min_val} {self.s3_object.max_val} {self.s3_object.max_val - self.s3_object.min_val + 1} {self.s3_object.max_val - self.s3_object.min_val + 1}
83                > DROP TABLE {staging_table}
84                > DROP CONNECTION {conn_name}
85                > DROP SECRET {secret_name}
86                """),
87            mz_service=state.mz_service,
88        )
class CopyToS3(materialize.zippy.framework.Action):
22class CopyToS3(Action):
23    """Performs a COPY TO S3 and records the resulting S3 object as a capability."""
24
25    @classmethod
26    def requires(cls) -> set[type[Capability]]:
27        return {BalancerdIsRunning, MzIsRunning, TableExists}
28
29    def __init__(self, capabilities: Capabilities) -> None:
30        self.table = random.choice(capabilities.get(TableExists))
31        super().__init__(capabilities)
32        self.s3_key = f"zippy/{self.seqno}"
33        self.s3_object = S3ObjectExists(
34            s3_key=self.s3_key,
35            min_val=self.table.watermarks.min,
36            max_val=self.table.watermarks.max,
37        )
38
39    def provides(self) -> list[Capability]:
40        return [self.s3_object]
41
42    def run(self, c: Composition, state: State) -> None:
43        conn_name = f"zippy_aws_conn_{self.seqno}"
44        secret_name = f"zippy_minio_secret_{self.seqno}"
45
46        c.testdrive(
47            dedent(f"""
48                > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}'
49                > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name})
50                > COPY (SELECT f1 FROM {self.table.name}) TO 's3://copytos3/{self.s3_key}' WITH (AWS CONNECTION = {conn_name}, FORMAT = 'csv')
51                > DROP CONNECTION {conn_name}
52                > DROP SECRET {secret_name}
53                """),
54            mz_service=state.mz_service,
55        )
56        self.s3_object.min_val = self.table.watermarks.min
57        self.s3_object.max_val = self.table.watermarks.max

Performs a COPY TO S3 and records the resulting S3 object as a capability.

CopyToS3(capabilities: materialize.zippy.framework.Capabilities)
29    def __init__(self, capabilities: Capabilities) -> None:
30        self.table = random.choice(capabilities.get(TableExists))
31        super().__init__(capabilities)
32        self.s3_key = f"zippy/{self.seqno}"
33        self.s3_object = S3ObjectExists(
34            s3_key=self.s3_key,
35            min_val=self.table.watermarks.min,
36            max_val=self.table.watermarks.max,
37        )

Construct a new action, possibly conditioning on the available capabilities.

@classmethod
def requires(cls) -> set[type[materialize.zippy.framework.Capability]]:
25    @classmethod
26    def requires(cls) -> set[type[Capability]]:
27        return {BalancerdIsRunning, MzIsRunning, TableExists}

Compute the capability classes that this action requires.

table
s3_key
s3_object
def provides(self) -> list[materialize.zippy.framework.Capability]:
39    def provides(self) -> list[Capability]:
40        return [self.s3_object]

Compute the capabilities that this action will make available.

def run( self, c: materialize.mzcompose.composition.Composition, state: materialize.zippy.framework.State) -> None:
42    def run(self, c: Composition, state: State) -> None:
43        conn_name = f"zippy_aws_conn_{self.seqno}"
44        secret_name = f"zippy_minio_secret_{self.seqno}"
45
46        c.testdrive(
47            dedent(f"""
48                > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}'
49                > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name})
50                > COPY (SELECT f1 FROM {self.table.name}) TO 's3://copytos3/{self.s3_key}' WITH (AWS CONNECTION = {conn_name}, FORMAT = 'csv')
51                > DROP CONNECTION {conn_name}
52                > DROP SECRET {secret_name}
53                """),
54            mz_service=state.mz_service,
55        )
56        self.s3_object.min_val = self.table.watermarks.min
57        self.s3_object.max_val = self.table.watermarks.max

Run this action on the provided composition.

class CopyFromS3(materialize.zippy.framework.Action):
60class CopyFromS3(Action):
61    """Reads a previously written S3 object back into a staging table and validates the row count."""
62
63    @classmethod
64    def requires(cls) -> set[type[Capability]]:
65        return {BalancerdIsRunning, MzIsRunning, S3ObjectExists}
66
67    def __init__(self, capabilities: Capabilities) -> None:
68        self.s3_object = random.choice(capabilities.get(S3ObjectExists))
69        super().__init__(capabilities)
70
71    def run(self, c: Composition, state: State) -> None:
72        conn_name = f"zippy_aws_conn_{self.seqno}"
73        secret_name = f"zippy_minio_secret_{self.seqno}"
74        staging_table = f"zippy_s3_staging_{self.seqno}"
75
76        c.testdrive(
77            dedent(f"""
78                > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}'
79                > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name})
80                > CREATE TABLE {staging_table} (f1 INTEGER)
81                > COPY INTO {staging_table} FROM 's3://copytos3/{self.s3_object.s3_key}' (FORMAT CSV, AWS CONNECTION = {conn_name})
82                > SELECT MIN(f1), MAX(f1), COUNT(f1), COUNT(DISTINCT f1) FROM {staging_table}
83                {self.s3_object.min_val} {self.s3_object.max_val} {self.s3_object.max_val - self.s3_object.min_val + 1} {self.s3_object.max_val - self.s3_object.min_val + 1}
84                > DROP TABLE {staging_table}
85                > DROP CONNECTION {conn_name}
86                > DROP SECRET {secret_name}
87                """),
88            mz_service=state.mz_service,
89        )

Reads a previously written S3 object back into a staging table and validates the row count.

CopyFromS3(capabilities: materialize.zippy.framework.Capabilities)
67    def __init__(self, capabilities: Capabilities) -> None:
68        self.s3_object = random.choice(capabilities.get(S3ObjectExists))
69        super().__init__(capabilities)

Construct a new action, possibly conditioning on the available capabilities.

@classmethod
def requires(cls) -> set[type[materialize.zippy.framework.Capability]]:
63    @classmethod
64    def requires(cls) -> set[type[Capability]]:
65        return {BalancerdIsRunning, MzIsRunning, S3ObjectExists}

Compute the capability classes that this action requires.

s3_object
def run( self, c: materialize.mzcompose.composition.Composition, state: materialize.zippy.framework.State) -> None:
71    def run(self, c: Composition, state: State) -> None:
72        conn_name = f"zippy_aws_conn_{self.seqno}"
73        secret_name = f"zippy_minio_secret_{self.seqno}"
74        staging_table = f"zippy_s3_staging_{self.seqno}"
75
76        c.testdrive(
77            dedent(f"""
78                > CREATE SECRET {secret_name} AS '${{testdrive.aws-secret-access-key}}'
79                > CREATE CONNECTION {conn_name} TO AWS (ENDPOINT '${{testdrive.aws-endpoint}}', REGION 'us-east-1', ACCESS KEY ID '${{testdrive.aws-access-key-id}}', SECRET ACCESS KEY SECRET {secret_name})
80                > CREATE TABLE {staging_table} (f1 INTEGER)
81                > COPY INTO {staging_table} FROM 's3://copytos3/{self.s3_object.s3_key}' (FORMAT CSV, AWS CONNECTION = {conn_name})
82                > SELECT MIN(f1), MAX(f1), COUNT(f1), COUNT(DISTINCT f1) FROM {staging_table}
83                {self.s3_object.min_val} {self.s3_object.max_val} {self.s3_object.max_val - self.s3_object.min_val + 1} {self.s3_object.max_val - self.s3_object.min_val + 1}
84                > DROP TABLE {staging_table}
85                > DROP CONNECTION {conn_name}
86                > DROP SECRET {secret_name}
87                """),
88            mz_service=state.mz_service,
89        )

Run this action on the provided composition.