Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
Go
at sl/gitmirror
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620package db
import ( "context" "database/sql" "errors" "fmt" "time")
type CacheEntry struct { ID string StorageKey string OwnerDID string RepoDID string Engine string CacheKey string CacheHash string Checksum string SizeBytes int64 RestoreCount int64 State string CreatedAt time.Time LastUsedAt time.Time}
const cacheEntryColumns = ` id, storage_key, owner_did, repo_did, engine, cache_key, cache_hash, checksum, size_bytes, restore_count, state, created_at, last_used_at`
func (d *DB) InsertCacheEntry(ctx context.Context, entry CacheEntry) error { _, err := d.ExecContext(ctx, ` insert into cache_entries ( id, storage_key, owner_did, repo_did, engine, cache_key, cache_hash, checksum, size_bytes, restore_count, state, created_at, last_used_at ) values (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, entry.ID, entry.StorageKey, entry.OwnerDID, entry.RepoDID, entry.Engine, entry.CacheKey, entry.CacheHash, entry.Checksum, entry.SizeBytes, entry.RestoreCount, entry.State, entry.CreatedAt.UnixNano(), entry.LastUsedAt.UnixNano(), ) return err}func (d *DB) InsertCacheEntryWithinQuota(ctx context.Context, entry CacheEntry, maxEntries int64) (bool, error) { result, err := d.ExecContext(ctx, ` insert into cache_entries ( id, storage_key, owner_did, repo_did, engine, cache_key, cache_hash, checksum, size_bytes, restore_count, state, created_at, last_used_at ) select ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ? where ? <= 0 or ( select count(*) from cache_entries where owner_did = ? and state in ('pending', 'ready', 'deleting') ) < ?`, entry.ID, entry.StorageKey, entry.OwnerDID, entry.RepoDID, entry.Engine, entry.CacheKey, entry.CacheHash, entry.Checksum, entry.SizeBytes, entry.RestoreCount, entry.State, entry.CreatedAt.UnixNano(), entry.LastUsedAt.UnixNano(), maxEntries, entry.OwnerDID, maxEntries, ) if err != nil { return false, err } inserted, err := result.RowsAffected() return inserted == 1, err}
func (d *DB) MarkCacheEntryReady(ctx context.Context, id string, sizeBytes, maxBytes int64, now time.Time) ([]CacheEntry, error) { return d.MarkCacheEntryReadyWithChecksum(ctx, id, sizeBytes, "", maxBytes, now)}
func (d *DB) MarkCacheEntryReadyWithChecksum(ctx context.Context, id string, sizeBytes int64, checksum string, maxBytes int64, now time.Time) ([]CacheEntry, error) { tx, err := d.BeginTx(ctx, nil) if err != nil { return nil, err } defer tx.Rollback() superseded, ready, err := markCacheEntryReady(ctx, tx, id, sizeBytes, checksum, maxBytes, now) if err != nil { return nil, err } if !ready { return nil, sql.ErrNoRows } if err := tx.Commit(); err != nil { return nil, err } return superseded, nil}
type cacheEntryTx interface { QueryRowContext(context.Context, string, ...any) *sql.Row QueryContext(context.Context, string, ...any) (*sql.Rows, error)}
func markCacheEntryReady(ctx context.Context, tx cacheEntryTx, id string, sizeBytes int64, checksum string, maxBytes int64, now time.Time) ([]CacheEntry, bool, error) { current, err := scanCacheEntry(tx.QueryRowContext(ctx, `select `+cacheEntryColumns+` from cache_entries where id = ?`, id)) if err == sql.ErrNoRows { return nil, false, nil } if err != nil { return nil, false, err } if current.State != "pending" { if current.State == "ready" { if current.SizeBytes != sizeBytes { return nil, false, fmt.Errorf("cache entry %q size changed from %d to %d", id, current.SizeBytes, sizeBytes) } return nil, true, nil } return nil, false, nil }
// A replacement frees its old generation. Exclude it while deciding how // much room is needed, then evict cold unrelated entries if necessary. rows, err := tx.QueryContext(ctx, ` update cache_entries set state = 'deleting' where repo_did = ? and engine = ? and cache_key = ? and cache_hash = ? and state = 'ready' returning `+cacheEntryColumns, current.RepoDID, current.Engine, current.CacheKey, current.CacheHash) if err != nil { return nil, false, err } var superseded []CacheEntry for rows.Next() { entry, err := scanCacheEntry(rows) if err != nil { rows.Close() return nil, false, err } superseded = append(superseded, *entry) } if err := rows.Close(); err != nil { return nil, false, err } if err := rows.Err(); err != nil { return nil, false, err }
if maxBytes > 0 { var used int64 if err := tx.QueryRowContext(ctx, ` select coalesce(sum(size_bytes), 0) from cache_entries where owner_did = ? and state in ('pending', 'ready', 'deleting') and id <> ? and not (repo_did = ? and engine = ? and cache_key = ? and cache_hash = ? and state = 'deleting')`, current.OwnerDID, current.ID, current.RepoDID, current.Engine, current.CacheKey, current.CacheHash).Scan(&used); err != nil { return nil, false, err } need := used + sizeBytes - maxBytes if need > 0 { rows, err := tx.QueryContext(ctx, ` select `+cacheEntryColumns+` from cache_entries where owner_did = ? and state = 'ready' and not (repo_did = ? and engine = ? and cache_key = ? and cache_hash = ?) order by last_used_at, created_at`, current.OwnerDID, current.RepoDID, current.Engine, current.CacheKey, current.CacheHash) if err != nil { return nil, false, err } var candidates []CacheEntry for rows.Next() { entry, err := scanCacheEntry(rows) if err != nil { rows.Close() return nil, false, err } candidates = append(candidates, *entry) } if err := rows.Close(); err != nil { return nil, false, err } if err := rows.Err(); err != nil { return nil, false, err } for _, entry := range candidates { if need <= 0 { break } claimed, err := tx.QueryContext(ctx, ` update cache_entries set state = 'deleting' where id = ? and state = 'ready' returning `+cacheEntryColumns, entry.ID) if err != nil { return nil, false, err } if claimed.Next() { deleted, scanErr := scanCacheEntry(claimed) claimed.Close() if scanErr != nil { return nil, false, scanErr } superseded = append(superseded, *deleted) need -= deleted.SizeBytes } else { claimed.Close() } } if need > 0 { return nil, false, nil } } }
updated, err := tx.QueryContext(ctx, ` update cache_entries set state = 'ready', checksum = ?, size_bytes = ?, last_used_at = ? where id = ? and state = 'pending' returning id`, checksum, sizeBytes, now.UnixNano(), id) if err != nil { return nil, false, err } ready := updated.Next() updated.Close() return superseded, ready, nil}
func (tx *EventBatchTx) MarkCacheEntryReady(ctx context.Context, id string, sizeBytes, maxBytes int64, now time.Time) ([]CacheEntry, bool, error) { return markCacheEntryReady(ctx, tx.tx, id, sizeBytes, "", maxBytes, now)}func (tx *EventBatchTx) MarkCacheEntryReadyWithChecksum(ctx context.Context, id string, sizeBytes int64, checksum string, maxBytes int64, now time.Time) ([]CacheEntry, bool, error) { return markCacheEntryReady(ctx, tx.tx, id, sizeBytes, checksum, maxBytes, now)}
// ReserveCacheEntryBytes atomically reserves the owner's currently available// byte budget in a pending row before its object is written.func (d *DB) ReserveCacheEntryBytes(ctx context.Context, id string, maxBytes int64) (int64, error) { if maxBytes <= 0 { return 0, nil } tx, err := d.BeginTx(ctx, nil) if err != nil { return 0, err } defer tx.Rollback() entry, err := scanCacheEntry(tx.QueryRowContext(ctx, `select `+cacheEntryColumns+` from cache_entries where id = ?`, id)) if err != nil { return 0, err } if entry.State != "pending" { return 0, sql.ErrNoRows } var used int64 if err := tx.QueryRowContext(ctx, ` select coalesce(sum(size_bytes), 0) from cache_entries where owner_did = ? and state in ('pending', 'ready', 'deleting') and id <> ? and not (repo_did = ? and engine = ? and cache_key = ? and cache_hash = ? and state in ('ready', 'deleting'))`, entry.OwnerDID, entry.ID, entry.RepoDID, entry.Engine, entry.CacheKey, entry.CacheHash).Scan(&used); err != nil { return 0, err } available := maxBytes - used if available <= 0 { return 0, fmt.Errorf("%w: owner %q has %d bytes in use", ErrCacheQuota, entry.OwnerDID, used) } if _, err := tx.ExecContext(ctx, `update cache_entries set size_bytes = ? where id = ? and state = 'pending'`, available, id); err != nil { return 0, err } if err := tx.Commit(); err != nil { return 0, err } return available, nil}
var ErrCacheQuota = errors.New("cache byte quota exceeded")
func (tx *EventBatchTx) TouchCacheEntry(ctx context.Context, id string, now time.Time) error { _, err := tx.tx.ExecContext(ctx, ` update cache_entries set last_used_at = max(last_used_at, ?) where id = ? and state = 'ready'`, now.UnixNano(), id) return err}
func (tx *EventBatchTx) DiscardCacheEntry(ctx context.Context, id string, pendingOnly bool) (*CacheEntry, error) { stateFilter := "" if pendingOnly { stateFilter = " and state in ('pending', 'deleting')" } current, err := scanCacheEntry(tx.tx.QueryRowContext(ctx, ` update cache_entries set state = 'deleting' where id = ?`+stateFilter+` returning `+cacheEntryColumns, id)) if err == sql.ErrNoRows { return nil, nil } return current, err}
func (d *DB) FindCacheEntry(ctx context.Context, repoDID, engine, key, hash string) (*CacheEntry, error) { return scanCacheEntry(d.QueryRowContext(ctx, ` select `+cacheEntryColumns+` from cache_entries where repo_did = ? and engine = ? and cache_key = ? and cache_hash = ? and state = 'ready' order by created_at desc limit 1`, repoDID, engine, key, hash))}
func (d *DB) FindFallbackCacheEntry(ctx context.Context, repoDID, engine, key, excludeHash string) (*CacheEntry, error) { return scanCacheEntry(d.QueryRowContext(ctx, ` select `+cacheEntryColumns+` from cache_entries where repo_did = ? and engine = ? and cache_key = ? and cache_hash <> ? and cache_hash <> '' and state = 'ready' order by created_at desc limit 1`, repoDID, engine, key, excludeHash))}func (d *DB) FindFallbackCacheEntryForIdentity(ctx context.Context, repoDID, engine, prefix, excludeHash string) (*CacheEntry, error) { return scanCacheEntry(d.QueryRowContext(ctx, ` select `+cacheEntryColumns+` from cache_entries where repo_did = ? and engine = ? and cache_key like ? || '%' and cache_hash <> ? and cache_hash <> '' and state = 'ready' order by created_at desc limit 1`, repoDID, engine, prefix, excludeHash))}
func (d *DB) BeginCacheRestore(ctx context.Context, id string) (bool, error) { result, err := d.ExecContext(ctx, `update cache_entries set restore_count = restore_count + 1 where id = ? and state = 'ready'`, id) if err != nil { return false, err } n, err := result.RowsAffected() return n == 1, err}
func (d *DB) EndCacheRestore(ctx context.Context, id string) error { _, err := d.ExecContext(ctx, `update cache_entries set restore_count = restore_count - 1 where id = ? and restore_count > 0`, id) return err}
func (d *DB) CacheEntryPinned(ctx context.Context, id string) (bool, error) { var count int64 err := d.QueryRowContext(ctx, `select restore_count from cache_entries where id = ?`, id).Scan(&count) if err == sql.ErrNoRows { return false, nil } return count > 0, err}
func (d *DB) TouchCacheEntry(ctx context.Context, id string, now time.Time) error { _, err := d.ExecContext(ctx, ` update cache_entries set last_used_at = max(last_used_at, ?) where id = ? and state = 'ready'`, now.UnixNano(), id) return err}
func (d *DB) ClaimPendingCacheEntry(ctx context.Context, id string) (bool, error) { result, err := d.ExecContext(ctx, ` update cache_entries set state = 'deleting' where id = ? and state = 'pending'`, id) if err != nil { return false, err } changed, err := result.RowsAffected() return changed == 1, err}
func (d *DB) ClaimCacheEntry(ctx context.Context, id, expectedState string, expectedLastUsed time.Time) (bool, error) { result, err := d.ExecContext(ctx, ` update cache_entries set state = 'deleting' where id = ? and state = ? and last_used_at = ? and restore_count = 0`, id, expectedState, expectedLastUsed.UnixNano()) if err != nil { return false, err } changed, err := result.RowsAffected() return changed == 1, err}
func (d *DB) RestoreCacheEntryState(ctx context.Context, id, state string) error { _, err := d.ExecContext(ctx, ` update cache_entries set state = ? where id = ? and state = 'deleting'`, state, id) return err}
func (d *DB) ExpiredCacheEntries(ctx context.Context, readyBefore, pendingBefore time.Time, limit int) ([]CacheEntry, error) { rows, err := d.QueryContext(ctx, ` select `+cacheEntryColumns+` from cache_entries where (state = 'ready' and last_used_at < ?) or (state in ('pending', 'deleting') and last_used_at < ?) order by last_used_at limit ?`, readyBefore.UnixNano(), pendingBefore.UnixNano(), limit, ) if err != nil { return nil, err } defer rows.Close()
var entries []CacheEntry for rows.Next() { entry, err := scanCacheEntry(rows) if err != nil { return nil, err } entries = append(entries, *entry) } return entries, rows.Err()}
// discardPendingCacheEntriesForLease fences uploads that can arrive after a lease dies.// Their tombstones let the cache pruner remove a late object without metadata.func (d *DB) DiscardPendingCacheEntriesForLease(leaseID string, now time.Time) ([]string, error) { tx, err := d.Begin() if err != nil { return nil, err } defer tx.Rollback()
rows, err := tx.Query(` select c.id, c.storage_key from cache_entries c join mill_cache_capabilities cap on cap.cache_id = c.id where cap.lease_id = ? and cap.action = 'save' and c.state = 'pending'`, leaseID) if err != nil { return nil, err } var keys []string var ids []string for rows.Next() { var id, key string if err := rows.Scan(&id, &key); err != nil { _ = rows.Close() return nil, err } ids = append(ids, id) keys = append(keys, key) } if err := rows.Err(); err != nil { _ = rows.Close() return nil, err } if err := rows.Close(); err != nil { return nil, err } for i, id := range ids { if _, err := tx.Exec(`update cache_entries set state = 'deleting' where id = ? and state = 'pending'`, id); err != nil { return nil, err } if _, err := tx.Exec(` insert into cache_object_deletions (storage_key, created_at) values (?, ?) on conflict(storage_key) do nothing`, keys[i], now.UnixNano()); err != nil { return nil, err } } if err := tx.Commit(); err != nil { return nil, err } return keys, nil}
func (d *DB) DeleteCacheEntry(ctx context.Context, id string) error { _, err := d.ExecContext(ctx, `delete from cache_entries where id = ?`, id) return err}
func (d *DB) DeleteCacheEntriesByStorageKey(ctx context.Context, key string) error { _, err := d.ExecContext(ctx, `delete from cache_entries where storage_key = ? and restore_count = 0`, key) return err}
type MillCacheCapability struct { Action string CacheID string StorageKey string}
func (d *DB) SaveMillCacheCapabilities(leaseID string, capabilities []MillCacheCapability) error { tx, err := d.Begin() if err != nil { return err } defer tx.Rollback() if _, err := tx.Exec(`delete from mill_cache_capabilities where lease_id = ?`, leaseID); err != nil { return err } for _, capability := range capabilities { if capability.Action != "restore" && capability.Action != "save" { return fmt.Errorf("cache capability has invalid action %q", capability.Action) } if capability.CacheID == "" || capability.StorageKey == "" { return fmt.Errorf("cache capability has empty identity") } if _, err := tx.Exec(` insert into mill_cache_capabilities ( lease_id, action, cache_id, storage_key ) values (?, ?, ?, ?)`, leaseID, capability.Action, capability.CacheID, capability.StorageKey, ); err != nil { return err } } return tx.Commit()}
func (tx *EventBatchTx) ConsumeMillCacheCapability(ctx context.Context, leaseID, action, cacheID string) (string, bool, error) { var storageKey string err := tx.tx.QueryRowContext(ctx, ` delete from mill_cache_capabilities where lease_id = ? and action = ? and cache_id = ? returning storage_key`, leaseID, action, cacheID, ).Scan(&storageKey) if err == sql.ErrNoRows { return "", false, nil } if err != nil { return "", false, err } return storageKey, true, nil}
func (tx *EventBatchTx) QueueCacheObjectDeletion(ctx context.Context, storageKey string, now time.Time) error { _, err := tx.tx.ExecContext(ctx, ` insert into cache_object_deletions (storage_key, created_at) values (?, ?) on conflict(storage_key) do nothing`, storageKey, now.UnixNano(), ) return err}
func (d *DB) PendingCacheObjectDeletions(ctx context.Context, limit int) ([]string, error) { rows, err := d.QueryContext(ctx, ` select storage_key from cache_object_deletions order by created_at limit ?`, limit) if err != nil { return nil, err } defer rows.Close() var keys []string for rows.Next() { var key string if err := rows.Scan(&key); err != nil { return nil, err } keys = append(keys, key) } return keys, rows.Err()}
func (d *DB) CacheEntriesByStorageKey(ctx context.Context, key string) ([]CacheEntry, error) { rows, err := d.QueryContext(ctx, `select `+cacheEntryColumns+` from cache_entries where storage_key = ?`, key) if err != nil { return nil, err } defer rows.Close() var entries []CacheEntry for rows.Next() { entry, err := scanCacheEntry(rows) if err != nil { return nil, err } entries = append(entries, *entry) } return entries, rows.Err()}
func (d *DB) CompleteCacheObjectDeletion(ctx context.Context, storageKey string) error { _, err := d.ExecContext(ctx, `delete from cache_object_deletions where storage_key = ?`, storageKey) return err}
type cacheEntryScanner interface { Scan(dest ...any) error}
func scanCacheEntry(row cacheEntryScanner) (*CacheEntry, error) { var entry CacheEntry var createdAt, lastUsedAt int64 if err := row.Scan( &entry.ID, &entry.StorageKey, &entry.OwnerDID, &entry.RepoDID, &entry.Engine, &entry.CacheKey, &entry.CacheHash, &entry.Checksum, &entry.SizeBytes, &entry.RestoreCount, &entry.State, &createdAt, &lastUsedAt, ); err != nil { return nil, err } entry.CreatedAt = time.Unix(0, createdAt) entry.LastUsedAt = time.Unix(0, lastUsedAt) return &entry, nil}