Skip to content

Commit cae843b

Browse files
mason-sharpclaude
andcommitted
Add hash version tracking for automatic merkle tree migration
Store a hash_version in ace_mtree_metadata so that mtree-update can detect trees built with an older hash algorithm and automatically recompute all leaf hashes. On first mtree-update after upgrade, the version mismatch marks all leaves dirty, triggering a full recompute within the normal update flow. This is crash-safe (version is updated in the same transaction as the hashes) and idempotent. Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
1 parent 8502e5d commit cae843b

3 files changed

Lines changed: 128 additions & 3 deletions

File tree

db/queries/queries.go

Lines changed: 64 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,11 @@ import (
2626
"github.com/pgedge/ace/pkg/types"
2727
)
2828

29+
// CurrentHashVersion is the version of the hash algorithm used by this build.
30+
// Increment when the SQL hash computation changes (e.g., switching from
31+
// whole-row ::text to per-column concat_ws with trim_scale).
32+
const CurrentHashVersion = 2
33+
2934
type DBQuerier interface {
3035
Exec(context.Context, string, ...interface{}) (pgconn.CommandTag, error)
3136
Query(context.Context, string, ...interface{}) (pgx.Rows, error)
@@ -1165,13 +1170,13 @@ func GetPkeyType(ctx context.Context, db DBQuerier, schema, table, pkey string)
11651170
return pkeyType, nil
11661171
}
11671172

1168-
func UpdateMetadata(ctx context.Context, db DBQuerier, schema, table string, totalRows int64, blockSize, numBlocks int, isComposite bool) error {
1173+
func UpdateMetadata(ctx context.Context, db DBQuerier, schema, table string, totalRows int64, blockSize, numBlocks int, isComposite bool, hashVersion int) error {
11691174
sql, err := RenderSQL(SQLTemplates.UpdateMetadata, nil)
11701175
if err != nil {
11711176
return fmt.Errorf("failed to render UpdateMetadata SQL: %w", err)
11721177
}
11731178

1174-
_, err = db.Exec(ctx, sql, schema, table, totalRows, blockSize, numBlocks, isComposite)
1179+
_, err = db.Exec(ctx, sql, schema, table, totalRows, blockSize, numBlocks, isComposite, hashVersion)
11751180
if err != nil {
11761181
return fmt.Errorf("query to update metadata for '%s.%s' failed: %w", schema, table, err)
11771182
}
@@ -2326,6 +2331,63 @@ func GetBlockSizeFromMetadata(ctx context.Context, db DBQuerier, schema, table s
23262331
return blockSize, nil
23272332
}
23282333

2334+
func EnsureHashVersionColumn(ctx context.Context, db DBQuerier) error {
2335+
sql, err := RenderSQL(SQLTemplates.EnsureHashVersionColumn, nil)
2336+
if err != nil {
2337+
return fmt.Errorf("failed to render EnsureHashVersionColumn SQL: %w", err)
2338+
}
2339+
2340+
_, err = db.Exec(ctx, sql)
2341+
if err != nil {
2342+
return fmt.Errorf("failed to add hash_version column to metadata table: %w", err)
2343+
}
2344+
return nil
2345+
}
2346+
2347+
func GetHashVersion(ctx context.Context, db DBQuerier, schema, table string) (int, error) {
2348+
sql, err := RenderSQL(SQLTemplates.GetHashVersion, nil)
2349+
if err != nil {
2350+
return 0, fmt.Errorf("failed to render GetHashVersion SQL: %w", err)
2351+
}
2352+
2353+
var version int
2354+
err = db.QueryRow(ctx, sql, schema, table).Scan(&version)
2355+
if err != nil {
2356+
return 0, fmt.Errorf("query to get hash version for '%s.%s' failed: %w", schema, table, err)
2357+
}
2358+
return version, nil
2359+
}
2360+
2361+
func MarkAllLeavesDirty(ctx context.Context, db DBQuerier, mtreeTable string) (int64, error) {
2362+
data := map[string]interface{}{
2363+
"MtreeTable": mtreeTable,
2364+
}
2365+
2366+
sql, err := RenderSQL(SQLTemplates.MarkAllLeavesDirty, data)
2367+
if err != nil {
2368+
return 0, fmt.Errorf("failed to render MarkAllLeavesDirty SQL: %w", err)
2369+
}
2370+
2371+
tag, err := db.Exec(ctx, sql)
2372+
if err != nil {
2373+
return 0, fmt.Errorf("query to mark all leaves dirty for '%s' failed: %w", mtreeTable, err)
2374+
}
2375+
return tag.RowsAffected(), nil
2376+
}
2377+
2378+
func UpdateHashVersion(ctx context.Context, db DBQuerier, schema, table string, version int) error {
2379+
sql, err := RenderSQL(SQLTemplates.UpdateHashVersion, nil)
2380+
if err != nil {
2381+
return fmt.Errorf("failed to render UpdateHashVersion SQL: %w", err)
2382+
}
2383+
2384+
_, err = db.Exec(ctx, sql, version, schema, table)
2385+
if err != nil {
2386+
return fmt.Errorf("query to update hash version for '%s.%s' failed: %w", schema, table, err)
2387+
}
2388+
return nil
2389+
}
2390+
23292391
func GetMaxNodeLevel(ctx context.Context, db DBQuerier, mtreeTable string) (int, error) {
23302392
data := map[string]interface{}{
23312393
"MtreeTable": mtreeTable,

db/queries/templates.go

Lines changed: 29 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -120,6 +120,10 @@ type Templates struct {
120120
RemoveTableFromCDCMetadata *template.Template
121121
GetSpockOriginLSNForNode *template.Template
122122
GetSpockSlotLSNForNode *template.Template
123+
EnsureHashVersionColumn *template.Template
124+
GetHashVersion *template.Template
125+
MarkAllLeavesDirty *template.Template
126+
UpdateHashVersion *template.Template
123127
}
124128

125129
var SQLTemplates = Templates{
@@ -132,6 +136,7 @@ var SQLTemplates = Templates{
132136
block_size int,
133137
num_blocks int,
134138
is_composite boolean NOT NULL DEFAULT false,
139+
hash_version int NOT NULL DEFAULT 2,
135140
last_updated timestamptz,
136141
PRIMARY KEY (schema_name, table_name)
137142
)`),
@@ -839,6 +844,7 @@ var SQLTemplates = Templates{
839844
block_size,
840845
num_blocks,
841846
is_composite,
847+
hash_version,
842848
last_updated
843849
)
844850
VALUES
@@ -849,6 +855,7 @@ var SQLTemplates = Templates{
849855
$4,
850856
$5,
851857
$6,
858+
$7,
852859
current_timestamp
853860
)
854861
ON CONFLICT (schema_name, table_name) DO
@@ -858,6 +865,7 @@ var SQLTemplates = Templates{
858865
block_size = EXCLUDED.block_size,
859866
num_blocks = EXCLUDED.num_blocks,
860867
is_composite = EXCLUDED.is_composite,
868+
hash_version = EXCLUDED.hash_version,
861869
last_updated = EXCLUDED.last_updated
862870
`)),
863871
DeleteMetadata: template.Must(template.New("deleteMetadata").Parse(`
@@ -1543,4 +1551,25 @@ var SQLTemplates = Templates{
15431551
ORDER BY rs.confirmed_flush_lsn DESC
15441552
LIMIT 1
15451553
`)),
1554+
EnsureHashVersionColumn: template.Must(template.New("ensureHashVersionColumn").Parse(`
1555+
ALTER TABLE spock.ace_mtree_metadata
1556+
ADD COLUMN IF NOT EXISTS hash_version int NOT NULL DEFAULT 1
1557+
`)),
1558+
GetHashVersion: template.Must(template.New("getHashVersion").Parse(`
1559+
SELECT COALESCE(
1560+
(SELECT hash_version FROM spock.ace_mtree_metadata
1561+
WHERE schema_name = $1 AND table_name = $2),
1562+
1
1563+
)
1564+
`)),
1565+
MarkAllLeavesDirty: template.Must(template.New("markAllLeavesDirty").Parse(`
1566+
UPDATE {{.MtreeTable}}
1567+
SET dirty = true
1568+
WHERE node_level = 0
1569+
`)),
1570+
UpdateHashVersion: template.Must(template.New("updateHashVersion").Parse(`
1571+
UPDATE spock.ace_mtree_metadata
1572+
SET hash_version = $1, last_updated = current_timestamp
1573+
WHERE schema_name = $2 AND table_name = $3
1574+
`)),
15461575
}

internal/consistency/mtree/merkle.go

Lines changed: 35 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1507,6 +1507,12 @@ func (m *MerkleTreeTask) UpdateMtree(skipAllChecks bool) (err error) {
15071507
return fmt.Errorf("error getting connection pool for node %s: %w", nodeInfo["Name"], err)
15081508
}
15091509

1510+
// Ensure hash_version column exists (schema migration for upgrades).
1511+
if err := queries.EnsureHashVersionColumn(m.Ctx, pool); err != nil {
1512+
pool.Close()
1513+
return fmt.Errorf("error migrating metadata schema on node %s: %w", nodeInfo["Name"], err)
1514+
}
1515+
15101516
blockSize, err = queries.GetBlockSizeFromMetadata(m.Ctx, pool, m.Schema, m.Table)
15111517
if err != nil {
15121518
pool.Close()
@@ -1556,12 +1562,34 @@ func (m *MerkleTreeTask) UpdateMtree(skipAllChecks bool) (err error) {
15561562
mtreeTableIdentifier := pgx.Identifier{aceSchema(), fmt.Sprintf("ace_mtree_%s_%s", m.Schema, m.Table)}
15571563
mtreeTableName := mtreeTableIdentifier.Sanitize()
15581564

1565+
// Check if stored hashes use an older algorithm and need full recomputation.
1566+
hashVersion, err := queries.GetHashVersion(m.Ctx, tx, m.Schema, m.Table)
1567+
if err != nil {
1568+
return fmt.Errorf("error getting hash version on node %s: %w", nodeInfo["Name"], err)
1569+
}
1570+
hashVersionUpgraded := false
1571+
if hashVersion < queries.CurrentHashVersion {
1572+
marked, err := queries.MarkAllLeavesDirty(m.Ctx, tx, mtreeTableName)
1573+
if err != nil {
1574+
return fmt.Errorf("error marking all leaves dirty for hash upgrade on node %s: %w", nodeInfo["Name"], err)
1575+
}
1576+
fmt.Printf("Hash algorithm upgraded (v%d -> v%d): marked %d blocks for recomputation on %s\n",
1577+
hashVersion, queries.CurrentHashVersion, marked, nodeInfo["Name"])
1578+
hashVersionUpgraded = true
1579+
}
1580+
15591581
blocksToUpdate, err := queries.GetDirtyAndNewBlocks(m.Ctx, tx, mtreeTableName, m.SimplePrimaryKey, m.Key)
15601582
if err != nil {
15611583
return fmt.Errorf("error getting dirty blocks on node %s: %w", nodeInfo["Name"], err)
15621584
}
15631585

15641586
if len(blocksToUpdate) == 0 {
1587+
if hashVersionUpgraded {
1588+
// No blocks exist yet, but still update the version marker.
1589+
if err := queries.UpdateHashVersion(m.Ctx, tx, m.Schema, m.Table, queries.CurrentHashVersion); err != nil {
1590+
return fmt.Errorf("error updating hash version on node %s: %w", nodeInfo["Name"], err)
1591+
}
1592+
}
15651593
fmt.Printf("No updates needed for %s\n", nodeInfo["Name"])
15661594
tx.Commit(m.Ctx)
15671595
continue
@@ -1642,6 +1670,12 @@ func (m *MerkleTreeTask) UpdateMtree(skipAllChecks bool) (err error) {
16421670
}
16431671
}
16441672

1673+
if hashVersionUpgraded {
1674+
if err := queries.UpdateHashVersion(m.Ctx, tx, m.Schema, m.Table, queries.CurrentHashVersion); err != nil {
1675+
return fmt.Errorf("error updating hash version on node %s: %w", nodeInfo["Name"], err)
1676+
}
1677+
}
1678+
16451679
if err := tx.Commit(m.Ctx); err != nil {
16461680
return fmt.Errorf("error committing transaction on node %s: %w", nodeInfo["Name"], err)
16471681
}
@@ -2440,7 +2474,7 @@ func (m *MerkleTreeTask) createMtreeObjects(tx pgx.Tx, totalRows int64, numBlock
24402474
return fmt.Errorf("failed to create metadata table: %w", err)
24412475
}
24422476

2443-
err = queries.UpdateMetadata(m.Ctx, tx, m.Schema, m.Table, totalRows, m.BlockSize, numBlocks, !m.SimplePrimaryKey)
2477+
err = queries.UpdateMetadata(m.Ctx, tx, m.Schema, m.Table, totalRows, m.BlockSize, numBlocks, !m.SimplePrimaryKey, queries.CurrentHashVersion)
24442478
if err != nil {
24452479
return fmt.Errorf("failed to update metadata: %w", err)
24462480
}

0 commit comments

Comments
 (0)