Skip to content
8 changes: 4 additions & 4 deletions dbos/dbos_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -251,7 +251,7 @@ func TestConfig(t *testing.T) {
require.NoError(t, err)
assert.True(t, exists, "dbos_migrations table should exist")

// Verify migration version is 9 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, and workflow_schedules migration)
// Verify migration version is 10 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, migration 9, and notifications primary key migration)
var version int64
var count int
err = sysDB.pool.QueryRow(dbCtx, "SELECT COUNT(*) FROM dbos.dbos_migrations").Scan(&count)
Expand All @@ -260,7 +260,7 @@ func TestConfig(t *testing.T) {

err = sysDB.pool.QueryRow(dbCtx, "SELECT version FROM dbos.dbos_migrations").Scan(&version)
require.NoError(t, err)
assert.Equal(t, int64(9), version, "migration version should be 9 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, and workflow_schedules migration)")
assert.Equal(t, int64(10), version, "migration version should be 10 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, migration 9, and notifications primary key migration)")

// Test manual shutdown and recreate
Shutdown(ctx, 1*time.Minute)
Expand Down Expand Up @@ -547,7 +547,7 @@ func TestCustomSystemDBSchema(t *testing.T) {
require.NoError(t, err)
assert.True(t, exists, "dbos_migrations table should exist in custom schema")

// Verify migration version is 9 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, and workflow_schedules migration)
// Verify migration version is 10 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, migration 9, and notifications primary key migration)
var version int64
var count int
err = sysDB.pool.QueryRow(dbCtx, fmt.Sprintf("SELECT COUNT(*) FROM %s.dbos_migrations", customSchema)).Scan(&count)
Expand All @@ -556,7 +556,7 @@ func TestCustomSystemDBSchema(t *testing.T) {

err = sysDB.pool.QueryRow(dbCtx, fmt.Sprintf("SELECT version FROM %s.dbos_migrations", customSchema)).Scan(&version)
require.NoError(t, err)
assert.Equal(t, int64(9), version, "migration version should be 9 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, and workflow_schedules migration)")
assert.Equal(t, int64(10), version, "migration version should be 10 (after initial migration, queue partition key migration, workflow status index migration, forked_from migration, step timestamps migration, workflow events history migration, owner_xid migration, parent_workflow_id migration, migration 9, and notifications primary key migration)")
})

// Test workflows for exercising Send/Recv and SetEvent/GetEvent
Expand Down
18 changes: 18 additions & 0 deletions dbos/migrations/10_add_notifications_pkey.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
-- Migration 10: Add primary key to notifications table
-- An earlier version of DBOS had a bug where this table was created without a primary key.
-- The initial migration has been changed to create a key, and this migration creates the key
-- for existing applications.

DO $$
BEGIN
IF NOT EXISTS (
SELECT 1 FROM pg_constraint c
JOIN pg_class cl ON c.conrelid = cl.oid
JOIN pg_namespace n ON cl.relnamespace = n.oid
WHERE n.nspname = '%s'
AND cl.relname = 'notifications'
AND c.contype = 'p'
) THEN
ALTER TABLE %s.notifications ADD PRIMARY KEY (message_uuid);
END IF;
END $$;
18 changes: 13 additions & 5 deletions dbos/system_database.go
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,9 @@ var migration8SQL string
//go:embed migrations/9_add_workflow_schedules.sql
var migration9SQL string

//go:embed migrations/10_add_notifications_pkey.sql
var migration10SQL string

type migrationFile struct {
version int64
sql string
Expand Down Expand Up @@ -224,9 +227,13 @@ func runMigrations(ctx context.Context, pool *pgxpool.Pool, schema string, isCoc

migration8SQLProcessed := fmt.Sprintf(migration8SQL, sanitizedSchema, sanitizedSchema)

migration9SQLProcessed := fmt.Sprintf(migration9SQL, sanitizedSchema)
migration9SQLProcessed := fmt.Sprintf(migration9SQL, sanitizedSchema)

migration10SQLProcessed := fmt.Sprintf(migration10SQL, schema, sanitizedSchema)

// Build migrations list with processed SQL
// NOTE: Migration 9 will be added from PR #261
// TODO: Add migration 9 once PR #261 is merged
migrations := []migrationFile{
{version: 1, sql: migration1SQLProcessed},
{version: 2, sql: migration2SQLProcessed},
Expand All @@ -237,6 +244,7 @@ func runMigrations(ctx context.Context, pool *pgxpool.Pool, schema string, isCoc
{version: 7, sql: migration7SQLProcessed},
{version: 8, sql: migration8SQLProcessed},
{version: 9, sql: migration9SQLProcessed},
{version: 10, sql: migration10SQLProcessed},
}

// Begin transaction for atomic migration execution
Expand Down Expand Up @@ -1162,7 +1170,7 @@ func (s *sysDB) resumeWorkflow(ctx context.Context, input resumeWorkflowDBInput)

// Set the workflow's status to ENQUEUED and clear its recovery attempts, set new deadline
updateStatusQuery := fmt.Sprintf(`UPDATE %s.workflow_status
SET status = $1, queue_name = $2, recovery_attempts = $3,
SET status = $1, queue_name = $2, recovery_attempts = $3,
workflow_deadline_epoch_ms = NULL, deduplication_id = NULL,
started_at_epoch_ms = NULL, updated_at = $4
WHERE workflow_uuid = $5`, pgx.Identifier{s.schema}.Sanitize())
Expand Down Expand Up @@ -2565,11 +2573,11 @@ func (s *sysDB) writeStream(ctx context.Context, input writeStreamDBInput) error
return s.pool.Exec(ctx, sql, args...)
}

checkClosedQuery := fmt.Sprintf(`SELECT 1 FROM %s.streams
checkClosedQuery := fmt.Sprintf(`SELECT 1 FROM %s.streams
WHERE workflow_uuid = $1 AND key = $2 AND value = $3 LIMIT 1`,
pgx.Identifier{s.schema}.Sanitize())

getOffsetQuery := fmt.Sprintf(`SELECT COALESCE(MAX("offset"), -1) + 1 FROM %s.streams
getOffsetQuery := fmt.Sprintf(`SELECT COALESCE(MAX("offset"), -1) + 1 FROM %s.streams
WHERE workflow_uuid = $1 AND key = $2`,
pgx.Identifier{s.schema}.Sanitize())

Expand Down Expand Up @@ -2604,7 +2612,7 @@ func (s *sysDB) writeStream(ctx context.Context, input writeStreamDBInput) error
// readStream reads stream entries starting from a given offset.
// Returns the entries, whether the stream is closed, and any error.
func (s *sysDB) readStream(ctx context.Context, input readStreamDBInput) ([]streamEntry, bool, error) {
query := fmt.Sprintf(`SELECT value, "offset" FROM %s.streams
query := fmt.Sprintf(`SELECT value, "offset" FROM %s.streams
WHERE workflow_uuid = $1 AND key = $2 AND "offset" >= $3
ORDER BY "offset" ASC`,
pgx.Identifier{s.schema}.Sanitize())
Expand Down
Loading