Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
3.7 kB · 114 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115package model
import ( "context" "errors" "fmt" "time"
"gorm.io/gorm")
type Teleport struct { URI string `json:"uri" gorm:"primaryKey;column:uri"` CID string `json:"cid" gorm:"column:cid"` StartsAt time.Time `json:"startsAt" gorm:"column:starts_at;index:idx_repo_starts,priority:2"` DurationSeconds *int64 `json:"durationSeconds" gorm:"column:duration_seconds"` ViewerCount int64 `json:"viewerCount" gorm:"column:viewer_count;default:0"` Teleport *[]byte `json:"teleport"` RepoDID string `json:"repoDID" gorm:"column:repo_did;index:idx_repo_starts,priority:1"` TargetDID string `json:"targetDID" gorm:"column:target_did;index:idx_target_did"` Denied bool `json:"denied" gorm:"column:denied;default:false"` Repo *Repo `json:"repo,omitempty" gorm:"foreignKey:DID;references:RepoDID"` Target *Repo `json:"target,omitempty" gorm:"foreignKey:DID;references:TargetDID"`}
// CreateTeleport upserts a teleport record. As with livestreams, an unchanged// record is now a no-op: re-indexing one used to reschedule its arrival// notification and re-publish it to the bus.func (m *DBModel) CreateTeleport(ctx context.Context, tp *Teleport) error { return createOrVerify(ctx, m, tp, map[string]any{"uri": tp.URI})}
func (m *DBModel) GetLatestTeleportForRepo(repoDID string) (*Teleport, error) { var teleport Teleport err := m.DB. Preload("Repo"). Preload("Target"). Where("repo_did = ?", repoDID). Order("starts_at DESC"). First(&teleport).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, fmt.Errorf("error retrieving latest teleport: %w", err) } return &teleport, nil}
func (m *DBModel) GetActiveTeleportsForRepo(repoDID string) ([]Teleport, error) { now := time.Now() var teleports []Teleport err := m.DB. Preload("Repo"). Preload("Target"). Where("repo_did = ?", repoDID). Where("denied = ?", false). Where("starts_at <= ?", now). Where("(duration_seconds IS NULL OR DATE_ADD(starts_at, INTERVAL duration_seconds SECOND) > ?)", now). Order("starts_at DESC"). Find(&teleports).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, fmt.Errorf("error retrieving active teleports: %w", err) } return teleports, nil}
func (m *DBModel) GetActiveTeleportsToRepo(targetDID string) ([]Teleport, error) { now := time.Now() var teleports []Teleport err := m.DB. Preload("Repo"). Preload("Target"). Where("target_did = ?", targetDID). Where("denied = ?", false). Where("starts_at <= ?", now). Where("(duration_seconds IS NULL OR datetime(starts_at, '+' || duration_seconds || ' seconds') > ?)", now). Order("starts_at DESC"). Find(&teleports).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, fmt.Errorf("error retrieving active teleports to repo: %w", err) } return teleports, nil}
func (m *DBModel) GetTeleportByURI(uri string) (*Teleport, error) { var teleport Teleport err := m.DB. Preload("Repo"). Preload("Target"). Where("uri = ?", uri). First(&teleport).Error if errors.Is(err, gorm.ErrRecordNotFound) { return nil, nil } if err != nil { return nil, fmt.Errorf("error retrieving teleport by uri: %w", err) } return &teleport, nil}
func (m *DBModel) DeleteTeleport(ctx context.Context, uri string) error { return m.DB.Where("uri = ?", uri).Delete(&Teleport{}).Error}
func (m *DBModel) DenyTeleport(ctx context.Context, uri string) error { return m.DB.Model(&Teleport{}).Where("uri = ?", uri).Update("denied", true).Error}