package store import ( "context" "database/sql" "errors" "time" ) // HookFailureLimit is chapter 23.4's twenty, after which the hook stops and the page says so. const HookFailureLimit = 20 // HookState is what the server remembers about one webhook, keyed by the url in .barerepo/config. type HookState struct { URL string Failures int Disabled bool LastError string LastAt time.Time } // HooksOf returns the delivery state of every hook this repository has ever had. func (db *DB) HooksOf(ctx context.Context, repo string) (map[string]HookState, error) { rows, err := db.QueryContext(ctx, `SELECT url, failures, disabled_at, last_error, last_at FROM webhooks WHERE repo = ?`, repo) if err != nil { return nil, err } defer rows.Close() out := map[string]HookState{} for rows.Next() { var h HookState var disabled, last sql.NullInt64 if err := rows.Scan(&h.URL, &h.Failures, &disabled, &h.LastError, &last); err != nil { return nil, err } h.Disabled = disabled.Valid if last.Valid { h.LastAt = time.Unix(last.Int64, 0) } out[h.URL] = h } return out, rows.Err() } // HookDelivered clears the count, because twenty consecutive failures means consecutive. func (db *DB) HookDelivered(ctx context.Context, repo, url string) error { return db.upsertHook(ctx, repo, url, false, "") } // HookFailed counts one failure and disables the hook at the limit. func (db *DB) HookFailed(ctx context.Context, repo, url, reason string) error { return db.upsertHook(ctx, repo, url, true, reason) } func (db *DB) upsertHook(ctx context.Context, repo, url string, failed bool, reason string) error { at := now().Unix() if !failed { res, err := db.ExecContext(ctx, `UPDATE webhooks SET failures = 0, disabled_at = NULL, last_error = '', last_at = ? WHERE repo = ? AND url = ?`, at, repo, url) if err != nil { return err } if n, _ := res.RowsAffected(); n > 0 { return nil } _, err = db.ExecContext(ctx, `INSERT INTO webhooks (repo, url, failures, last_error, last_at) VALUES (?, ?, 0, '', ?)`, repo, url, at) return err } res, err := db.ExecContext(ctx, `UPDATE webhooks SET failures = failures + 1, last_error = ?, last_at = ?, disabled_at = CASE WHEN failures + 1 >= ? THEN ? ELSE disabled_at END WHERE repo = ? AND url = ?`, reason, at, HookFailureLimit, at, repo, url) if err != nil { return err } if n, _ := res.RowsAffected(); n > 0 { return nil } _, err = db.ExecContext(ctx, `INSERT INTO webhooks (repo, url, failures, last_error, last_at) VALUES (?, ?, 1, ?, ?)`, repo, url, reason, at) return err } // ForgetHook drops the state, so an operator who fixed the receiver can start it again. func (db *DB) ForgetHook(ctx context.Context, repo, url string) error { _, err := db.ExecContext(ctx, `DELETE FROM webhooks WHERE repo = ? AND url = ?`, repo, url) return err } // WebhookCursor is the last event dispatched, so a restart neither repeats nor skips. func (db *DB) WebhookCursor(ctx context.Context) (int64, error) { var id int64 err := db.QueryRowContext(ctx, `SELECT event_id FROM webhook_cursor WHERE id = 1`).Scan(&id) if errors.Is(err, sql.ErrNoRows) { return 0, nil } return id, err } func (db *DB) SetWebhookCursor(ctx context.Context, id int64) error { res, err := db.ExecContext(ctx, `UPDATE webhook_cursor SET event_id = ? WHERE id = 1`, id) if err != nil { return err } if n, _ := res.RowsAffected(); n > 0 { return nil } _, err = db.ExecContext(ctx, `INSERT INTO webhook_cursor (id, event_id) VALUES (1, ?)`, id) return err } // StartWebhooksHere puts the cursor at the newest event, so a first start sends no backlog. func (db *DB) StartWebhooksHere(ctx context.Context) error { if _, err := db.WebhookCursor(ctx); err != nil { return err } var id sql.NullInt64 if err := db.QueryRowContext(ctx, `SELECT MAX(id) FROM events`).Scan(&id); err != nil { return err } var have int if err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM webhook_cursor WHERE id = 1`).Scan(&have); err != nil { return err } if have > 0 { return nil } return db.SetWebhookCursor(ctx, id.Int64) } // EventsAfter returns the events a hook has not seen, oldest first, which is the order they happened. func (db *DB) EventsAfter(ctx context.Context, after int64, limit int) ([]Event, error) { rows, err := db.QueryContext(ctx, `SELECT id, kind, actor, repo, ref, number, title, detail, created_at FROM events WHERE id > ? ORDER BY id ASC LIMIT ?`, after, limit) if err != nil { return nil, err } defer rows.Close() return scanEvents(rows) }