mirror of
https://tangled.org/evan.jarrett.net/at-container-registry
synced 2026-09-25 11:44:16 +00:00
some request crawl relay fixes
This commit is contained in:
@@ -1,7 +1,6 @@
|
||||
package admin
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"net/http"
|
||||
"net/url"
|
||||
@@ -85,71 +84,87 @@ func (ui *AdminUI) handleRelayStatus(w http.ResponseWriter, r *http.Request) {
|
||||
ui.renderTemplate(w, "partials/relay_status.html", view)
|
||||
}
|
||||
|
||||
// handleRelayCrawl requests crawl from a single relay.
|
||||
func (ui *AdminUI) handleRelayCrawl(w http.ResponseWriter, r *http.Request) {
|
||||
if err := r.ParseForm(); err != nil {
|
||||
setFlash(w, r, "error", "Invalid form data")
|
||||
http.Redirect(w, r, "/admin#relays", http.StatusFound)
|
||||
return
|
||||
}
|
||||
|
||||
relayURL := r.FormValue("url")
|
||||
if relayURL == "" {
|
||||
setFlash(w, r, "error", "Missing relay URL")
|
||||
http.Redirect(w, r, "/admin#relays", http.StatusFound)
|
||||
return
|
||||
}
|
||||
|
||||
if err := atproto.RequestCrawl(relayURL, ui.config.PublicURL); err != nil {
|
||||
slog.Warn("Failed to request crawl from relay", "relay", relayURL, "error", err)
|
||||
setFlash(w, r, "error", "Crawl request failed: "+err.Error())
|
||||
} else {
|
||||
slog.Info("Crawl requested via admin panel", "relay", relayURL)
|
||||
setFlash(w, r, "success", "Crawl requested from "+relayURL)
|
||||
}
|
||||
|
||||
http.Redirect(w, r, "/admin#relays", http.StatusFound)
|
||||
// RelayCrawlResultView is the data for a single relay crawl result row.
|
||||
type RelayCrawlResultView struct {
|
||||
Name string
|
||||
URL string
|
||||
Success bool
|
||||
Error string
|
||||
}
|
||||
|
||||
// handleRelayCrawlAll requests crawl from all known relays.
|
||||
// handleRelayCrawl requests crawl from a single relay and returns an HTMX partial.
|
||||
func (ui *AdminUI) handleRelayCrawl(w http.ResponseWriter, r *http.Request) {
|
||||
relayURL := r.URL.Query().Get("url")
|
||||
relayName := r.URL.Query().Get("name")
|
||||
if relayURL == "" {
|
||||
http.Error(w, "Missing relay URL", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
err := atproto.RequestCrawl(relayURL, ui.config.PublicURL)
|
||||
|
||||
view := RelayCrawlResultView{
|
||||
Name: relayName,
|
||||
URL: relayURL,
|
||||
Success: err == nil,
|
||||
}
|
||||
if err != nil {
|
||||
view.Error = err.Error()
|
||||
slog.Warn("Failed to request crawl from relay", "relay", relayURL, "error", err)
|
||||
} else {
|
||||
slog.Info("Crawl requested via admin panel", "relay", relayURL)
|
||||
}
|
||||
|
||||
ui.renderTemplate(w, "partials/relay_crawl_result.html", view)
|
||||
}
|
||||
|
||||
// handleRelayCrawlAll requests crawl from all known relays and returns HTMX partials.
|
||||
func (ui *AdminUI) handleRelayCrawlAll(w http.ResponseWriter, r *http.Request) {
|
||||
relays := atproto.KnownRelays
|
||||
|
||||
var (
|
||||
wg sync.WaitGroup
|
||||
mu sync.Mutex
|
||||
successes int
|
||||
failures int
|
||||
)
|
||||
type result struct {
|
||||
relay atproto.KnownRelay
|
||||
err error
|
||||
}
|
||||
|
||||
for _, relay := range relays {
|
||||
results := make([]result, len(relays))
|
||||
var wg sync.WaitGroup
|
||||
|
||||
for i, relay := range relays {
|
||||
wg.Add(1)
|
||||
go func(relay atproto.KnownRelay) {
|
||||
go func(i int, relay atproto.KnownRelay) {
|
||||
defer wg.Done()
|
||||
if err := atproto.RequestCrawl(relay.URL, ui.config.PublicURL); err != nil {
|
||||
slog.Warn("Failed to request crawl", "relay", relay.Name, "error", err)
|
||||
mu.Lock()
|
||||
failures++
|
||||
mu.Unlock()
|
||||
} else {
|
||||
mu.Lock()
|
||||
successes++
|
||||
mu.Unlock()
|
||||
}
|
||||
}(relay)
|
||||
err := atproto.RequestCrawl(relay.URL, ui.config.PublicURL)
|
||||
results[i] = result{relay: relay, err: err}
|
||||
}(i, relay)
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
|
||||
session := getSessionFromContext(r.Context())
|
||||
successes := 0
|
||||
for _, res := range results {
|
||||
if res.err == nil {
|
||||
successes++
|
||||
}
|
||||
}
|
||||
slog.Info("Crawl all requested via admin panel",
|
||||
"successes", successes, "failures", failures, "by", session.DID)
|
||||
"successes", successes, "failures", len(relays)-successes, "by", session.DID)
|
||||
|
||||
if failures == 0 {
|
||||
setFlash(w, r, "success", fmt.Sprintf("Crawl requested from all %d relays", successes))
|
||||
} else {
|
||||
setFlash(w, r, "warning", fmt.Sprintf("Crawl: %d succeeded, %d failed", successes, failures))
|
||||
var views []RelayCrawlResultView
|
||||
for _, res := range results {
|
||||
v := RelayCrawlResultView{
|
||||
Name: res.relay.Name,
|
||||
URL: res.relay.URL,
|
||||
Success: res.err == nil,
|
||||
}
|
||||
if res.err != nil {
|
||||
v.Error = res.err.Error()
|
||||
}
|
||||
views = append(views, v)
|
||||
}
|
||||
|
||||
http.Redirect(w, r, "/admin#relays", http.StatusFound)
|
||||
ui.renderTemplate(w, "partials/relay_crawl_results.html", struct {
|
||||
Results []RelayCrawlResultView
|
||||
}{Results: views})
|
||||
}
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
{{define "partials/relay_crawl_result.html"}}
|
||||
<tr hx-get="/admin/api/relay/status?url={{.URL}}&name={{.Name}}"
|
||||
hx-trigger="load delay:10s"
|
||||
hx-swap="outerHTML">
|
||||
<td>
|
||||
{{if .Success}}
|
||||
<span class="badge badge-success badge-sm gap-1">
|
||||
{{ icon "check-circle" "size-3" }}
|
||||
Sent
|
||||
</span>
|
||||
{{else}}
|
||||
<span class="badge badge-error badge-sm gap-1">
|
||||
{{ icon "alert-circle" "size-3" }}
|
||||
Failed
|
||||
</span>
|
||||
{{end}}
|
||||
</td>
|
||||
<td>
|
||||
<div>
|
||||
<strong>{{.Name}}</strong><br>
|
||||
<code class="text-xs text-base-content/50 font-mono">{{.URL}}</code>
|
||||
</div>
|
||||
</td>
|
||||
<td>
|
||||
<span class="text-base-content/30 text-sm">-</span>
|
||||
</td>
|
||||
<td>
|
||||
{{if .Success}}
|
||||
<span class="text-sm text-info flex items-center gap-1">
|
||||
<span class="loading loading-spinner loading-xs"></span>
|
||||
Crawl requested, refreshing...
|
||||
</span>
|
||||
{{else}}
|
||||
<span class="text-sm text-error">{{.Error}}</span>
|
||||
{{end}}
|
||||
</td>
|
||||
<td class="text-right">
|
||||
<button class="btn btn-ghost btn-sm gap-1"
|
||||
hx-get="/admin/api/relay/status?url={{.URL}}&name={{.Name}}"
|
||||
hx-target="closest tr"
|
||||
hx-swap="outerHTML"
|
||||
title="Refresh Status">
|
||||
{{ icon "refresh-ccw" "size-4" }}
|
||||
Refresh
|
||||
</button>
|
||||
</td>
|
||||
</tr>
|
||||
{{end}}
|
||||
@@ -0,0 +1,50 @@
|
||||
{{define "partials/relay_crawl_results.html"}}
|
||||
{{range .Results}}
|
||||
<tr hx-get="/admin/api/relay/status?url={{.URL}}&name={{.Name}}"
|
||||
hx-trigger="load delay:10s"
|
||||
hx-swap="outerHTML">
|
||||
<td>
|
||||
{{if .Success}}
|
||||
<span class="badge badge-success badge-sm gap-1">
|
||||
{{ icon "check-circle" "size-3" }}
|
||||
Sent
|
||||
</span>
|
||||
{{else}}
|
||||
<span class="badge badge-error badge-sm gap-1">
|
||||
{{ icon "alert-circle" "size-3" }}
|
||||
Failed
|
||||
</span>
|
||||
{{end}}
|
||||
</td>
|
||||
<td>
|
||||
<div>
|
||||
<strong>{{.Name}}</strong><br>
|
||||
<code class="text-xs text-base-content/50 font-mono">{{.URL}}</code>
|
||||
</div>
|
||||
</td>
|
||||
<td>
|
||||
<span class="text-base-content/30 text-sm">-</span>
|
||||
</td>
|
||||
<td>
|
||||
{{if .Success}}
|
||||
<span class="text-sm text-info flex items-center gap-1">
|
||||
<span class="loading loading-spinner loading-xs"></span>
|
||||
Crawl requested, refreshing...
|
||||
</span>
|
||||
{{else}}
|
||||
<span class="text-sm text-error">{{.Error}}</span>
|
||||
{{end}}
|
||||
</td>
|
||||
<td class="text-right">
|
||||
<button class="btn btn-ghost btn-sm gap-1"
|
||||
hx-get="/admin/api/relay/status?url={{.URL}}&name={{.Name}}"
|
||||
hx-target="closest tr"
|
||||
hx-swap="outerHTML"
|
||||
title="Refresh Status">
|
||||
{{ icon "refresh-ccw" "size-4" }}
|
||||
Refresh
|
||||
</button>
|
||||
</td>
|
||||
</tr>
|
||||
{{end}}
|
||||
{{end}}
|
||||
@@ -48,13 +48,14 @@
|
||||
</td>
|
||||
<td class="text-right">
|
||||
{{if .Online}}
|
||||
<form action="/admin/relays/crawl" method="POST" class="inline">
|
||||
<input type="hidden" name="url" value="{{.URL}}">
|
||||
<button type="submit" class="btn btn-ghost btn-sm gap-1" title="Request Crawl">
|
||||
{{ icon "refresh-ccw" "size-4" }}
|
||||
Crawl
|
||||
</button>
|
||||
</form>
|
||||
<button class="btn btn-ghost btn-sm gap-1"
|
||||
hx-post="/admin/relays/crawl?url={{.URL}}&name={{.Name}}"
|
||||
hx-target="closest tr"
|
||||
hx-swap="outerHTML"
|
||||
title="Request Crawl">
|
||||
{{ icon "refresh-ccw" "size-4" }}
|
||||
Request Crawl
|
||||
</button>
|
||||
{{end}}
|
||||
</td>
|
||||
</tr>
|
||||
|
||||
@@ -1,12 +1,24 @@
|
||||
{{define "partials/tab_relays.html"}}
|
||||
<div class="flex justify-between items-center min-h-12 mb-6">
|
||||
<h1 class="text-2xl font-bold">Relays</h1>
|
||||
<form action="/admin/relays/crawl-all" method="POST">
|
||||
<button type="submit" class="btn btn-primary gap-2">
|
||||
{{ icon "refresh-ccw" "size-4" }}
|
||||
Crawl All
|
||||
</button>
|
||||
</form>
|
||||
<button class="btn btn-primary gap-2"
|
||||
hx-post="/admin/relays/crawl-all"
|
||||
hx-target="#relay-tbody"
|
||||
hx-swap="innerHTML"
|
||||
hx-indicator="#crawl-loading">
|
||||
{{ icon "refresh-ccw" "size-4" }}
|
||||
Request Crawl All
|
||||
</button>
|
||||
</div>
|
||||
|
||||
<div id="crawl-loading" class="htmx-indicator mb-4">
|
||||
<div class="flex items-center gap-3 p-4 bg-base-200 rounded-lg">
|
||||
<span class="loading loading-spinner loading-md text-primary"></span>
|
||||
<div>
|
||||
<p class="font-medium">Requesting crawl from all relays...</p>
|
||||
<p class="text-sm text-base-content/50">This may take a few seconds.</p>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<div class="card bg-base-100 shadow-sm">
|
||||
@@ -21,7 +33,7 @@
|
||||
<th class="text-right">Actions</th>
|
||||
</tr>
|
||||
</thead>
|
||||
<tbody>
|
||||
<tbody id="relay-tbody">
|
||||
{{range .Relays}}
|
||||
<tr hx-get="/admin/api/relay/status?url={{.URL}}&name={{.Name}}"
|
||||
hx-trigger="load"
|
||||
|
||||
@@ -387,6 +387,13 @@ func (b *EventBroadcaster) Subscribe(conn *websocket.Conn, cursor int64, userAge
|
||||
|
||||
slog.Info("New firehose subscriber", "remote", conn.RemoteAddr(), "cursor", cursor, "currentSeq", currentSeq, "userAgent", userAgent)
|
||||
|
||||
// Send #account event first so relays know this DID is active.
|
||||
// This is written directly to the WebSocket before backfill/handleSubscriber
|
||||
// start, so there are no concurrent writers at this point.
|
||||
if err := b.sendAccountEvent(conn); err != nil {
|
||||
slog.Warn("Failed to send account event to subscriber", "error", err)
|
||||
}
|
||||
|
||||
// Handle cursor-based backfill:
|
||||
// - cursor < 0: No backfill, stream new events only
|
||||
// - cursor >= 0: Backfill events from cursor onwards
|
||||
@@ -414,6 +421,41 @@ func (b *EventBroadcaster) Subscribe(conn *websocket.Conn, cursor int64, userAge
|
||||
return sub
|
||||
}
|
||||
|
||||
// sendAccountEvent writes an #account event directly to a WebSocket connection,
|
||||
// signaling that this DID is active on this host. This is critical for relays
|
||||
// that previously saw the DID deactivated on a different PDS.
|
||||
func (b *EventBroadcaster) sendAccountEvent(conn *websocket.Conn) error {
|
||||
header := events.EventHeader{
|
||||
Op: events.EvtKindMessage,
|
||||
MsgType: "#account",
|
||||
}
|
||||
|
||||
wc, err := conn.NextWriter(websocket.BinaryMessage)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to get websocket writer: %w", err)
|
||||
}
|
||||
|
||||
if err := header.MarshalCBOR(wc); err != nil {
|
||||
wc.Close()
|
||||
return fmt.Errorf("failed to write account event header: %w", err)
|
||||
}
|
||||
|
||||
acctEvt := &atproto.SyncSubscribeRepos_Account{
|
||||
Active: true,
|
||||
Did: b.holdDID,
|
||||
Seq: 0, // Not sequenced in the commit stream
|
||||
Time: time.Now().Format(time.RFC3339),
|
||||
}
|
||||
|
||||
var obj lexutil.CBOR = acctEvt
|
||||
if err := obj.MarshalCBOR(wc); err != nil {
|
||||
wc.Close()
|
||||
return fmt.Errorf("failed to write account event body: %w", err)
|
||||
}
|
||||
|
||||
return wc.Close()
|
||||
}
|
||||
|
||||
// Unsubscribe removes a WebSocket subscriber
|
||||
func (b *EventBroadcaster) Unsubscribe(sub *Subscriber) {
|
||||
b.mu.Lock()
|
||||
|
||||
Reference in New Issue
Block a user