1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
|
package main
import (
"database/sql"
"fmt"
"time"
_ "modernc.org/sqlite"
)
// The database holds the film list as *generations*: every row of every data
// table carries the id of the generation it belongs to, and the generation
// table says which one is served:
//
// building being written by an import; the server never looks at it
// live what the server answers from (at most one)
// previous the one before, kept for `mediathek generations rollback`
//
// Deleting a generation row deletes everything in it: every data table
// references generation(id) ON DELETE CASCADE. The full-text indexes are
// virtual tables, which a cascade cannot reach, so they are kept in step by
// AFTER DELETE triggers on the tables they index — the cascade fires those.
// See mediathek(7), "GENERATIONS".
const schema = `
CREATE TABLE IF NOT EXISTS generation (
id INTEGER PRIMARY KEY,
state TEXT NOT NULL CHECK (state IN ('building','live','previous')),
list_created INTEGER NOT NULL, -- film list header time, unix seconds UTC
list_hash TEXT NOT NULL,
started_at INTEGER NOT NULL,
finished_at INTEGER,
rows INTEGER NOT NULL DEFAULT 0, -- film list rows parsed
row_errors INTEGER NOT NULL DEFAULT 0,
dropped INTEGER NOT NULL DEFAULT 0, -- live streams, junk
videos INTEGER NOT NULL DEFAULT 0, -- after dedup
aliases INTEGER NOT NULL DEFAULT 0,
channels INTEGER NOT NULL DEFAULT 0,
senders INTEGER NOT NULL DEFAULT 0
);
CREATE UNIQUE INDEX IF NOT EXISTS generation_one_live ON generation(state) WHERE state = 'live';
CREATE UNIQUE INDEX IF NOT EXISTS generation_one_previous ON generation(state) WHERE state = 'previous';
-- A rowid table (not WITHOUT ROWID): the rowid is what video_fts indexes.
CREATE TABLE IF NOT EXISTS video (
gen INTEGER NOT NULL REFERENCES generation(id) ON DELETE CASCADE,
uuid TEXT NOT NULL,
channel TEXT NOT NULL,
family TEXT NOT NULL,
sender TEXT NOT NULL,
topic TEXT NOT NULL,
title TEXT NOT NULL,
description TEXT NOT NULL,
website TEXT NOT NULL,
geo TEXT NOT NULL,
published INTEGER NOT NULL,
duration INTEGER NOT NULL,
kind INTEGER NOT NULL,
files TEXT NOT NULL, -- see encodeFiles
subtitle TEXT NOT NULL,
tags TEXT NOT NULL, -- comma-separated
dkey INTEGER NOT NULL, -- episode dedup key, see dedupKey
score INTEGER NOT NULL -- higher wins dedup, see dedupScore
);
CREATE UNIQUE INDEX IF NOT EXISTS video_uuid ON video(gen, uuid);
CREATE INDEX IF NOT EXISTS video_channel ON video(gen, channel, published DESC);
CREATE INDEX IF NOT EXISTS video_family ON video(gen, family, published DESC);
CREATE INDEX IF NOT EXISTS video_published ON video(gen, published DESC);
-- Covers the dedup window (PARTITION BY dkey ORDER BY score DESC, uuid) so
-- it walks the index in order instead of sorting and reading every row.
DROP INDEX IF EXISTS video_dkey;
CREATE INDEX IF NOT EXISTS video_dedup ON video(gen, dkey, score DESC, uuid);
-- uuids of episode duplicates that dedup dropped, pointing at the kept one.
CREATE TABLE IF NOT EXISTS alias (
gen INTEGER NOT NULL REFERENCES generation(id) ON DELETE CASCADE,
uuid TEXT NOT NULL,
target TEXT NOT NULL,
PRIMARY KEY (gen, uuid)
) WITHOUT ROWID;
-- id is the numeric PeerTube id (numericID of the name), stable across
-- generations and therefore not the rowid; channel_fts indexes the rowid.
CREATE TABLE IF NOT EXISTS channel (
gen INTEGER NOT NULL REFERENCES generation(id) ON DELETE CASCADE,
name TEXT NOT NULL,
id INTEGER NOT NULL,
family TEXT NOT NULL,
display TEXT NOT NULL,
senders TEXT NOT NULL,
videos INTEGER NOT NULL,
latest INTEGER NOT NULL,
UNIQUE (gen, name),
UNIQUE (gen, id)
);
CREATE TABLE IF NOT EXISTS account (
gen INTEGER NOT NULL REFERENCES generation(id) ON DELETE CASCADE,
name TEXT NOT NULL,
id INTEGER NOT NULL,
display TEXT NOT NULL,
channels INTEGER NOT NULL,
videos INTEGER NOT NULL,
PRIMARY KEY (gen, name)
) WITHOUT ROWID;
-- Full-text search. Contentless (the text lives in video/channel only) with
-- contentless_delete so rows can be removed by rowid alone. A contentless
-- table returns NULL for its columns, so a search is restricted to one
-- generation by joining back on rowid and filtering video.gen/channel.gen.
--
-- Rows are indexed by the import once dedup is done (never on insert), so a
-- dedup loser is deleted before it was ever indexed; deleting an unknown
-- rowid from a contentless_delete table is a no-op.
CREATE VIRTUAL TABLE IF NOT EXISTS video_fts USING fts5(
title, topic,
content='', contentless_delete=1,
tokenize='unicode61 remove_diacritics 2');
CREATE TRIGGER IF NOT EXISTS video_fts_delete AFTER DELETE ON video BEGIN
DELETE FROM video_fts WHERE rowid = old.rowid;
END;
-- Thumbnail cache, filled by the server (see thumbs.go). Deliberately NOT
-- part of a generation: it is keyed by the video uuid, which is stable
-- across imports, so a lookup done once serves every later generation. The
-- importer prunes rows whose uuid is in no generation any more. url = ''
-- means "looked, nothing found"; failed marks a lookup that errored, which
-- is retried sooner.
CREATE TABLE IF NOT EXISTS thumb (
uuid TEXT PRIMARY KEY,
url TEXT NOT NULL,
checked_at INTEGER NOT NULL,
failed INTEGER NOT NULL DEFAULT 0
) WITHOUT ROWID;
CREATE VIRTUAL TABLE IF NOT EXISTS channel_fts USING fts5(
display,
content='', contentless_delete=1,
tokenize='unicode61 remove_diacritics 2');
CREATE TRIGGER IF NOT EXISTS channel_fts_delete AFTER DELETE ON channel BEGIN
DELETE FROM channel_fts WHERE rowid = old.rowid;
END;
`
// Busy timeouts. The server only reads, and in WAL mode a reader never waits
// for the writer, so a read that blocks this long means something is wrong.
// The importer waits out whatever else holds the write lock.
const (
readBusyTimeout = 5000
writeBusyTimeout = 30000
)
// dsn builds the SQLite connection string. The pragmas MUST travel in the DSN
// rather than a post-open PRAGMA statement: busy_timeout and foreign_keys are
// per-connection state and db.Exec only ever runs on one pooled connection,
// which would leave every other connection at SQLite's defaults ("fail
// immediately when locked", "ignore ON DELETE CASCADE"). The modernc driver
// applies _pragma to every connection it opens. The spelling _busy_timeout=
// is silently ignored by this driver.
func dsn(path string, busyTimeout int) string {
return fmt.Sprintf(
"file:%s?_pragma=journal_mode(WAL)&_pragma=busy_timeout(%d)&_pragma=foreign_keys(1)",
path, busyTimeout,
)
}
// openDB opens (creating if necessary) the database and applies the schema.
func openDB(path string, busyTimeout int) (*sql.DB, error) {
db, err := sql.Open("sqlite", dsn(path, busyTimeout))
if err != nil {
return nil, fmt.Errorf("open db %q: %w", path, err)
}
if _, err := db.Exec(schema); err != nil {
db.Close()
return nil, fmt.Errorf("apply schema: %w", err)
}
return db, nil
}
// openWriteDB opens the server's own writer, used only for the thumbnail
// cache. One connection makes it a mutex: our writes queue in database/sql
// instead of fighting over SQLite's write lock. It is separate from the read
// pool so that a write waiting out an import (the long busy timeout) never
// holds a connection a read needs.
func openWriteDB(path string) (*sql.DB, error) {
db, err := sql.Open("sqlite", dsn(path, writeBusyTimeout))
if err != nil {
return nil, fmt.Errorf("open write db %q: %w", path, err)
}
db.SetMaxOpenConns(1)
db.SetMaxIdleConns(1)
return db, nil
}
// Generation is a row of the generation table.
type Generation struct {
ID int64
State string
ListCreated time.Time
ListHash string
StartedAt time.Time
FinishedAt time.Time
Rows int
RowErrors int
Dropped int
Videos int
Aliases int
Channels int
Senders int
}
const genCols = `id, state, list_created, list_hash, started_at, coalesce(finished_at,0), rows, row_errors, dropped, videos, aliases, channels, senders`
func scanGen(sc interface{ Scan(...any) error }) (Generation, error) {
var g Generation
var lc, sa, fa int64
err := sc.Scan(&g.ID, &g.State, &lc, &g.ListHash, &sa, &fa, &g.Rows, &g.RowErrors, &g.Dropped, &g.Videos, &g.Aliases, &g.Channels, &g.Senders)
g.ListCreated, g.StartedAt = time.Unix(lc, 0).UTC(), time.Unix(sa, 0).UTC()
if fa != 0 {
g.FinishedAt = time.Unix(fa, 0).UTC()
}
return g, err
}
func listGenerations(db *sql.DB) ([]Generation, error) {
rows, err := db.Query(`SELECT ` + genCols + ` FROM generation ORDER BY id`)
if err != nil {
return nil, err
}
defer rows.Close()
var gs []Generation
for rows.Next() {
g, err := scanGen(rows)
if err != nil {
return nil, err
}
gs = append(gs, g)
}
return gs, rows.Err()
}
// liveGeneration returns the live generation, or ok=false if there is none.
func liveGeneration(db *sql.DB) (Generation, bool, error) {
g, err := scanGen(db.QueryRow(`SELECT ` + genCols + ` FROM generation WHERE state = 'live'`))
if err == sql.ErrNoRows {
return g, false, nil
}
return g, err == nil, err
}
// deleteGeneration removes a generation and, by cascade, all its rows.
func deleteGeneration(db *sql.DB, id int64) error {
_, err := db.Exec(`DELETE FROM generation WHERE id = ?`, id)
return err
}
// rollback swaps live and previous.
func rollback(db *sql.DB) error {
tx, err := db.Begin()
if err != nil {
return err
}
defer tx.Rollback()
var live, prev int64
if err := tx.QueryRow(`SELECT id FROM generation WHERE state='live'`).Scan(&live); err != nil {
return fmt.Errorf("no live generation: %w", err)
}
if err := tx.QueryRow(`SELECT id FROM generation WHERE state='previous'`).Scan(&prev); err != nil {
return fmt.Errorf("no previous generation: %w", err)
}
// Three steps because the one-live/one-previous indexes are checked per
// row, so a single swapping UPDATE would trip them halfway.
steps := []struct {
q string
id int64
}{
{`UPDATE generation SET state='building' WHERE id=?`, live}, // park
{`UPDATE generation SET state='live' WHERE id=?`, prev},
{`UPDATE generation SET state='previous' WHERE id=?`, live},
}
for _, s := range steps {
if _, err := tx.Exec(s.q, s.id); err != nil {
return err
}
}
return tx.Commit()
}
|