| 1 | package sessioncatalog |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "database/sql" |
| 6 | "errors" |
| 7 | ) |
| 8 | |
| 9 | // Discovery has already read bounded source metadata. Compare every projected |
| 10 | // field, not just size/mtime: this shortcut never certifies transcript content. |
| 11 | func sameMetadataProjection(existing, incoming SessionRecord) bool { |
| 12 | return existing.MissingSince == 0 && sameSessionIndexInput(existing, incoming) && |
| 13 | existing.RecoveryCopy == incoming.RecoveryCopy && |
| 14 | existing.RecoveryGroupID == incoming.RecoveryGroupID && |
| 15 | existing.RecoveryRole == incoming.RecoveryRole && |
| 16 | existing.RecoveryCanonical == incoming.RecoveryCanonical && |
| 17 | existing.LogicalTopicID == incoming.LogicalTopicID && |
| 18 | existing.OrdinaryVisible == incoming.OrdinaryVisible && |
| 19 | existing.LogFormat == incoming.LogFormat && existing.HeadCount == incoming.HeadCount && |
| 20 | existing.SelectedHeadID == incoming.SelectedHeadID |
| 21 | } |
| 22 | |
| 23 | func refreshUnchangedMetadata(ctx context.Context, tx *sql.Tx, incoming SessionRecord, pathKey string, generation int64) (bool, error) { |
| 24 | // Recheck at the publication boundary; a pre-transaction observation alone |
| 25 | // cannot authorize skipping a concurrent metadata update. |
| 26 | current, err := scanSession(tx.QueryRowContext(ctx, `SELECT `+sessionSelectColumns+` FROM catalog_sessions WHERE path_key=?`, pathKey)) |
| 27 | if errors.Is(err, sql.ErrNoRows) { |
| 28 | return false, nil |
| 29 | } |
| 30 | if err != nil || !sameMetadataProjection(current, incoming) { |
| 31 | return false, err |
| 32 | } |
| 33 | _, err = tx.ExecContext(ctx, `UPDATE catalog_sessions SET seen_generation=MAX(seen_generation,?) WHERE path_key=?`, generation, pathKey) |
| 34 | return err == nil, err |
| 35 | } |
| 36 |