﻿# MySql Scripts

<!-- Version variant: sqlpersistence_6; default: [/persistence/sql/mysql-scripts.md](/persistence/sql/mysql-scripts.md) -->


Scripts and SQL used when interacting with a [MySql](https://www.mysql.com/) database.

## Build Time

Scripts are created at build time and are executed as part of a deployment or decommissioning of an endpoint.

### Outbox

#### Create Table

<!-- snippet: MySql_OutboxCreateSql -->

```sql
set @tableNameQuoted = concat('`', @tablePrefix, 'OutboxData`');
set @tableNameNonQuoted = concat(@tablePrefix, 'OutboxData');

set @createTable =  concat('
    create table if not exists ', @tableNameQuoted, '(
        MessageId nvarchar(200) not null,
        Dispatched bit not null default 0,
        DispatchedAt datetime,
        PersistenceVersion varchar(23) not null,
        Operations json not null,
        primary key (MessageId)
    ) default charset=ascii;
');
prepare script from @createTable;
execute script;
deallocate prepare script;

select count(*)
into @exist
from information_schema.statistics
where
    table_schema = database() and
    index_name = 'Index_DispatchedAt' and
    table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('create index Index_DispatchedAt on ', @tableNameQuoted, '(DispatchedAt)'), 'select \'Index Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;

select count(*)
into @exist
from information_schema.statistics
where
    table_schema = database() and
    index_name = 'Index_Dispatched' and
    table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('create index Index_Dispatched on ', @tableNameQuoted, '(Dispatched)'), 'select \'Index Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

#### Drop Table

<!-- snippet: MySql_OutboxDropSql -->

```sql
set @tableName = concat('`', @tablePrefix, 'OutboxData`');

set @dropTable = concat('drop table if exists ', @tableName);
prepare script from @dropTable;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

### Saga

For a Saga with the following structure

<!-- snippet: CreationScriptSaga -->

```cs
[SqlSaga(transitionalCorrelationProperty: nameof(OrderSagaData.OrderId))]
public class OrderSaga : Saga<OrderSaga.OrderSagaData>,
    IAmStartedByMessages<StartSaga>
{
    protected override void ConfigureHowToFindSaga(SagaPropertyMapper<OrderSagaData> mapper)
    {
        mapper.ConfigureMapping<StartSaga>(msg => msg.OrderNumber).ToSaga(saga => saga.OrderNumber);
    }

    public class OrderSagaData :
        ContainSagaData
    {
        public int OrderNumber { get; set; }
        public Guid OrderId { get; set; }
    }
```

<!-- endsnippet -->

#### Create Table

<!-- snippet: MySql_SagaCreateSql -->

```sql
/* TableNameVariable */

set @tableNameQuoted = concat('`', @tablePrefix, 'OrderSaga`');
set @tableNameNonQuoted = concat(@tablePrefix, 'OrderSaga');

/* Initialize */

drop procedure if exists sqlpersistence_raiseerror;
create procedure sqlpersistence_raiseerror(message varchar(256))
begin
signal sqlstate
    'ERROR'
set
    message_text = message,
    mysql_errno = '45000';
end;

/* CreateTable */

set @createTable = concat('
    create table if not exists ', @tableNameQuoted, '(
        Id varchar(38) not null,
        Metadata json not null,
        Data json not null,
        PersistenceVersion varchar(23) not null,
        SagaTypeVersion varchar(23) not null,
        Concurrency int not null,
        primary key (Id)
    ) default charset=ascii;
');
prepare script from @createTable;
execute script;
deallocate prepare script;

/* AddProperty OrderNumber */

select count(*)
into @exist
from information_schema.columns
where table_schema = database() and
      column_name = 'Correlation_OrderNumber' and
      table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('alter table ', @tableNameQuoted, ' add column Correlation_OrderNumber bigint(20)'), 'select \'Column Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;

/* VerifyColumnType Int */

set @column_type_OrderNumber = (
  select concat(column_type,' character set ', character_set_name)
  from information_schema.columns
  where
    table_schema = database() and
    table_name = @tableNameNonQuoted and
    column_name = 'Correlation_OrderNumber'
);

set @query = IF(
    @column_type_OrderNumber <> 'bigint(20)',
    'call sqlpersistence_raiseerror(concat(\'Incorrect data type for Correlation_OrderNumber. Expected bigint(20) got \', @column_type_OrderNumber, \'.\'));',
    'select \'Column Type OK\' status');

prepare script from @query;
execute script;
deallocate prepare script;

/* WriteCreateIndex OrderNumber */

select count(*)
into @exist
from information_schema.statistics
where
    table_schema = database() and
    index_name = 'Index_Correlation_OrderNumber' and
    table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('create unique index Index_Correlation_OrderNumber on ', @tableNameQuoted, '(Correlation_OrderNumber)'), 'select \'Index Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;

/* AddProperty OrderId */

select count(*)
into @exist
from information_schema.columns
where table_schema = database() and
      column_name = 'Correlation_OrderId' and
      table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('alter table ', @tableNameQuoted, ' add column Correlation_OrderId varchar(38) character set ascii'), 'select \'Column Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;

/* VerifyColumnType Guid */

set @column_type_OrderId = (
  select concat(column_type,' character set ', character_set_name)
  from information_schema.columns
  where
    table_schema = database() and
    table_name = @tableNameNonQuoted and
    column_name = 'Correlation_OrderId'
);

set @query = IF(
    @column_type_OrderId <> 'varchar(38) character set ascii',
    'call sqlpersistence_raiseerror(concat(\'Incorrect data type for Correlation_OrderId. Expected varchar(38) character set ascii got \', @column_type_OrderId, \'.\'));',
    'select \'Column Type OK\' status');

prepare script from @query;
execute script;
deallocate prepare script;

/* CreateIndex OrderId */

select count(*)
into @exist
from information_schema.statistics
where
    table_schema = database() and
    index_name = 'Index_Correlation_OrderId' and
    table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('create unique index Index_Correlation_OrderId on ', @tableNameQuoted, '(Correlation_OrderId)'), 'select \'Index Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;

/* PurgeObsoleteIndex */

select concat('drop index ', index_name, ' on ', @tableNameQuoted, ';')
from information_schema.statistics
where
    table_schema = database() and
    table_name = @tableNameNonQuoted and
    index_name like 'Index_Correlation_%' and
    index_name <> 'Index_Correlation_OrderNumber' and
    index_name <> 'Index_Correlation_OrderId' and
    table_schema = database()
into @dropIndexQuery;
select if (
    @dropIndexQuery is not null,
    @dropIndexQuery,
    'select ''no index to delete'';')
    into @dropIndexQuery;

prepare script from @dropIndexQuery;
execute script;
deallocate prepare script;

/* PurgeObsoleteProperties */

select concat('alter table ', table_name, ' drop column ', column_name, ';')
from information_schema.columns
where
    table_schema = database() and
    table_name = @tableNameNonQuoted and
    column_name like 'Correlation_%' and
    column_name <> 'Correlation_OrderNumber' and
    column_name <> 'Correlation_OrderId'
into @dropPropertiesQuery;

select if (
    @dropPropertiesQuery is not null,
    @dropPropertiesQuery,
    'select ''no property to delete'';')
    into @dropPropertiesQuery;

prepare script from @dropPropertiesQuery;
execute script;
deallocate prepare script;

/* CompleteSagaScript */
```

<!-- endsnippet -->

#### Drop Table

<!-- snippet: MySql_SagaDropSql -->

```sql
/* TableNameVariable */

set @tableNameQuoted = concat('`', @tablePrefix, 'OrderSaga`');
set @tableNameNonQuoted = concat(@tablePrefix, 'OrderSaga');

/* DropTable */

set @dropTable = concat('drop table if exists ', @tableNameQuoted);
prepare script from @dropTable;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

### Subscription

#### Create Table

<!-- snippet: MySql_SubscriptionCreateSql -->

```sql
set @tableName = concat('`', @tablePrefix, 'SubscriptionData`');

set @createTable = concat('
    create table if not exists ', @tableName, '(
        Subscriber nvarchar(200) not null,
        Endpoint nvarchar(200),
        MessageType nvarchar(200) not null,
        PersistenceVersion varchar(23) not null,
        primary key clustered (Subscriber, MessageType)
    ) default charset=ascii;
');
prepare script from @createTable;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

#### Drop Table

<!-- snippet: MySql_SubscriptionDropSql -->

```sql
set @tableName = concat('`', @tablePrefix, 'SubscriptionData`');

set @dropTable = concat('drop table if exists ', @tableName);
prepare script from @dropTable;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

### Timeout

#### Create Table

<!-- snippet: MySql_TimeoutCreateSql -->

```sql
set @tableNameQuoted = concat('`', @tablePrefix, 'TimeoutData`');
set @tableNameNonQuoted = concat(@tablePrefix, 'TimeoutData');

set @createTable = concat('
    create table if not exists ', @tableNameQuoted, '(
        Id varchar(38) not null,
        Destination nvarchar(200),
        SagaId varchar(38),
        State longblob,
        Time datetime,
        Headers json not null,
        PersistenceVersion varchar(23) not null,
        primary key (Id)
    ) default charset=ascii;
');
prepare script from @createTable;
execute script;
deallocate prepare script;

select count(*)
into @exist
from information_schema.statistics
where
    table_schema = database() and
    index_name = 'Index_SagaId' and
    table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('create index Index_SagaId on ', @tableNameQuoted, '(SagaId)'), 'select \'Index Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;

select count(*)
into @exist
from information_schema.statistics
where
    table_schema = database() and
    index_name = 'Index_Time' and
    table_name = @tableNameNonQuoted;

set @query = IF(
    @exist <= 0,
    concat('create index Index_Time on ', @tableNameQuoted, '(Time)'), 'select \'Index Exists\' status');

prepare script from @query;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

#### Drop Table

<!-- snippet: MySql_TimeoutDropSql -->

```sql
set @tableName = concat('`', @tablePrefix, 'TimeoutData`');

set @dropTable = concat('drop table if exists ', @tableName, '');
prepare script from @dropTable;
execute script;
deallocate prepare script;
```

<!-- endsnippet -->

## Run Time

SQL used at runtime to query and update data.

### Outbox

Used at intervals to cleanup old outbox records.

<!-- snippet: MySql_OutboxCleanupSql -->

```sql
delete from `EndpointNameOutboxData`
where Dispatched = true and
      DispatchedAt < @DispatchedBefore
limit @BatchSize
```

<!-- endsnippet -->

#### Get

Used by `IOutboxStorage.SetAsDispatched`.

<!-- snippet: MySql_OutboxGetSql -->

```sql
select
    Dispatched,
    Operations
from `EndpointNameOutboxData`
where MessageId = @MessageId
```

<!-- endsnippet -->

#### SetAsDispatched

Used by `IOutboxStorage.SetAsDispatched`.

<!-- snippet: MySql_OutboxSetAsDispatchedSql -->

```sql
update `EndpointNameOutboxData`
set
    Dispatched = 1,
    DispatchedAt = @DispatchedAt,
    Operations = '[]'
where MessageId = @MessageId
```

<!-- endsnippet -->

#### Store

Used by `IOutboxStorage.Store`.

##### Optimistic (default) mode

<!-- snippet: MySql_OutboxOptimisticStoreSql -->

```sql
insert into `EndpointNameOutboxData`
(
    MessageId,
    Operations,
    PersistenceVersion
)
values
(
    @MessageId,
    @Operations,
    @PersistenceVersion
)
```

<!-- endsnippet -->

##### Pessimistic mode

<!-- snippet: MySql_OutboxPessimisticBeginSql -->

```sql
insert into `EndpointNameOutboxData`
(
    MessageId,
    Operations,
    PersistenceVersion
)
values
(
    @MessageId,
    '[]',
    @PersistenceVersion
)
```

<!-- endsnippet -->

<!-- snippet: MySql_OutboxPessimisticCompleteSql -->

```sql
update `EndpointNameOutboxData`
set
    Operations = @Operations
where MessageId = @MessageId
```

<!-- endsnippet -->


### Saga

#### Complete

Used by `ISagaPersister.Complete`.

<!-- snippet: MySql_SagaCompleteSql -->

```sql
delete from EndpointName_SagaName
where Id = @Id and Concurrency = @Concurrency
```

<!-- endsnippet -->

#### Save

Used by `ISagaPersister.Save`.

<!-- snippet: MySql_SagaSaveSql -->

```sql
insert into EndpointName_SagaName
(
    Id,
    Metadata,
    Data,
    PersistenceVersion,
    SagaTypeVersion,
    Concurrency,
    Correlation_CorrelationProperty,
    Correlation_TransitionalCorrelationProperty
)
values
(
    @Id,
    @Metadata,
    @Data,
    @PersistenceVersion,
    @SagaTypeVersion,
    1,
    @CorrelationId,
    @TransitionalCorrelationId
)
```

<!-- endsnippet -->

#### GetByProperty

Used by `ISagaPersister.Get(propertyName...)`.

<!-- snippet: MySql_SagaGetByPropertySql -->

```sql
select
    Id,
    SagaTypeVersion,
    Concurrency,
    Metadata,
    Data
from EndpointName_SagaName
where Correlation_PropertyName = @propertyValue
for update
```

<!-- endsnippet -->

#### GetBySagaId

Used by `ISagaPersister.Get(sagaId...)`.

<!-- snippet: MySql_SagaGetBySagaIdSql -->

```sql
select
    Id,
    SagaTypeVersion,
    Concurrency,
    Metadata,
    Data
from EndpointName_SagaName
where Id = @Id
for update
```

<!-- endsnippet -->

#### Update

Used by `ISagaPersister.Update`.

<!-- snippet: MySql_SagaUpdateSql -->

```sql
update EndpointName_SagaName
set
    Data = @Data,
    PersistenceVersion = @PersistenceVersion,
    SagaTypeVersion = @SagaTypeVersion,
    Concurrency = @Concurrency + 1,
    Correlation_TransitionalCorrelationProperty = @TransitionalCorrelationId
where
    Id = @Id and Concurrency = @Concurrency
```

<!-- endsnippet -->

#### Select used by Saga Finder

<!-- snippet: MySql_SagaSelectSql -->

```sql
select
    Id,
    SagaTypeVersion,
    Concurrency,
    Metadata,
    Data
from EndpointName_SagaName
where 1 = 1
for update
```

<!-- endsnippet -->


### Subscription

#### GetSubscribers

Used by `ISubscriptionStorage.GetSubscriberAddressesForMessage`.

<!-- snippet: MySql_SubscriptionGetSubscribersSql -->

```sql
select distinct Subscriber, Endpoint
from `EndpointNameSubscriptionData`
where MessageType in (@type0)
```

<!-- endsnippet -->

#### Subscribe

Used by `ISubscriptionStorage.Subscribe`.

<!-- snippet: MySql_SubscriptionSubscribeSql -->

```sql
insert into `EndpointNameSubscriptionData`
(
    Subscriber,
    MessageType,
    Endpoint,
    PersistenceVersion
)
values
(
    @Subscriber,
    @MessageType,
    @Endpoint,
    @PersistenceVersion
)
on duplicate key update
    Endpoint = coalesce(@Endpoint, Endpoint),
    PersistenceVersion = @PersistenceVersion
```

<!-- endsnippet -->

#### Unsubscribe

Used by `ISubscriptionStorage.Unsubscribe`.

<!-- snippet: MySql_SubscriptionUnsubscribeSql -->

```sql
delete from `EndpointNameSubscriptionData`
where
    Subscriber = @Subscriber and
    MessageType = @MessageType
```

<!-- endsnippet -->

### Timeout

#### Peek

Used by `IPersistTimeouts.Peek`.

<!-- snippet: MySql_TimeoutPeekSql -->

```sql
select
    Destination,
    SagaId,
    State,
    Time,
    Headers
from `EndpointNameTimeoutData`
where Id = @Id
```

<!-- endsnippet -->

#### Add

Used by `IPersistTimeouts.Add`.

<!-- snippet: MySql_TimeoutAddSql -->

```sql
insert into `EndpointNameTimeoutData`
(
    Id,
    Destination,
    SagaId,
    State,
    Time,
    Headers,
    PersistenceVersion
)
values
(
    @Id,
    @Destination,
    @SagaId,
    @State,
    @Time,
    @Headers,
    @PersistenceVersion
)
```

<!-- endsnippet -->

#### GetNextChunk

Used by `IQueryTimeouts.GetNextChunk`.

<!-- snippet: MySql_TimeoutNextSql -->

```sql
select Time from `EndpointNameTimeoutData`
where Time > @EndTime
order by Time
limit 1
```

<!-- endsnippet -->

<!-- snippet: MySql_TimeoutRangeSql -->

```sql
select Id, Time
from `EndpointNameTimeoutData`
where Time > @StartTime and Time <= @EndTime
```

<!-- endsnippet -->

#### TryRemove

Used by `IPersistTimeouts.TryRemove`.

<!-- snippet: MySql_TimeoutRemoveByIdSql -->

```sql
delete from `EndpointNameTimeoutData`
where Id = @Id;
```

<!-- endsnippet -->

#### RemoveTimeoutBy

Used by `IPersistTimeouts.RemoveTimeoutBy`.

<!-- snippet: MySql_TimeoutRemoveBySagaIdSql -->

```sql
delete from `EndpointNameTimeoutData`
where SagaId = @SagaId
```

<!-- endsnippet -->
