Merge pull request #40 from JMR-dev/ticket-14-publish-issues
#14 Decrypt/format reports and publish as GitHub issues (de-duped)
This commit was merged in pull request #40.
This commit is contained in:
@@ -0,0 +1,204 @@
|
||||
package publish
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/ingest"
|
||||
)
|
||||
|
||||
// MaxIssueBodyRunes is GitHub's hard cap on an issue body (ADR #6 §2.1 note:
|
||||
// "GitHub issue bodies are capped at 65,536 characters"). formatIssue keeps the
|
||||
// rendered body at or below this, truncating the free-text report if needed.
|
||||
const MaxIssueBodyRunes = 65536
|
||||
|
||||
// maxTitleRunes keeps titles well under GitHub's 256-char limit.
|
||||
const maxTitleRunes = 200
|
||||
|
||||
// truncationMarker is appended to a report whose text had to be cut to fit.
|
||||
const truncationMarker = "\n\n[… truncated: report exceeded GitHub's 65,536-character issue body limit …]"
|
||||
|
||||
// formatIssue renders a decrypted (scrubbed) report into a GitHub issue title and
|
||||
// body. plaintext is the JSON produced at ingest (an ingest.Report shape) after
|
||||
// PII scrubbing; if it does not parse as that shape, the raw text is embedded
|
||||
// verbatim so nothing is silently dropped.
|
||||
//
|
||||
// # Markdown-injection safety
|
||||
//
|
||||
// The report is user-supplied (scrubbed, but not trusted): its free text is
|
||||
// wrapped in a fenced code block (with a fence long enough to survive any backtick
|
||||
// run in the content) and metadata values are wrapped in inline code, so @mentions
|
||||
// and links in a report cannot notify users or render as active Markdown in the
|
||||
// issue. The body is finally clamped to MaxIssueBodyRunes runes as a backstop.
|
||||
func formatIssue(id string, plaintext []byte) (title, body string) {
|
||||
var rep ingest.Report
|
||||
parsed := json.Unmarshal(plaintext, &rep) == nil
|
||||
|
||||
title = buildTitle(rep, parsed)
|
||||
body = buildBody(id, rep, plaintext, parsed)
|
||||
return title, body
|
||||
}
|
||||
|
||||
// buildTitle produces a short, human-readable issue title.
|
||||
func buildTitle(rep ingest.Report, parsed bool) string {
|
||||
if !parsed {
|
||||
return "Automated bug report"
|
||||
}
|
||||
app := oneLine(rep.AppVersion)
|
||||
platform := oneLine(rep.Platform)
|
||||
title := "Bug report"
|
||||
if app != "" {
|
||||
title += ": " + app
|
||||
}
|
||||
if platform != "" {
|
||||
title += " on " + platform
|
||||
}
|
||||
return truncateRunes(title, maxTitleRunes)
|
||||
}
|
||||
|
||||
// buildBody assembles the Markdown body, fitting the report text within the rune
|
||||
// budget left after the fixed sections.
|
||||
func buildBody(id string, rep ingest.Report, plaintext []byte, parsed bool) string {
|
||||
var head strings.Builder
|
||||
head.WriteString("This issue was filed automatically by the LibreMail bug-report ingest pipeline ")
|
||||
head.WriteString("(opt-in debug reports from the app). PII was best-effort scrubbed before storage; ")
|
||||
head.WriteString("see `docs/privacy.md`.\n\n")
|
||||
head.WriteString("**Report ID:** `")
|
||||
head.WriteString(sanitizeCode(id))
|
||||
head.WriteString("`\n\n")
|
||||
|
||||
// The free-text report to embed: the report field when the payload parsed,
|
||||
// otherwise the whole decrypted blob (so a schema drift never loses content).
|
||||
reportText := string(plaintext)
|
||||
if parsed {
|
||||
head.WriteString(metadataTable(rep))
|
||||
head.WriteString("\n### Report\n\n")
|
||||
reportText = rep.Report
|
||||
} else {
|
||||
head.WriteString("_Payload did not match the expected report schema; showing the decrypted content verbatim._\n\n")
|
||||
head.WriteString("### Decrypted content\n\n")
|
||||
}
|
||||
|
||||
const footer = "\n\n---\nLabels: `bug-report`, `automated`, `needs-triage`.\n"
|
||||
|
||||
fence := fenceFor(reportText)
|
||||
// Budget for the report text = cap − (everything else) − the two fence lines
|
||||
// and their newlines.
|
||||
overhead := utf8.RuneCountInString(head.String()) +
|
||||
utf8.RuneCountInString(footer) +
|
||||
2*utf8.RuneCountInString(fence) +
|
||||
2 // the two newlines wrapping the fenced content
|
||||
budget := MaxIssueBodyRunes - overhead
|
||||
|
||||
if utf8.RuneCountInString(reportText) > budget {
|
||||
keep := budget - utf8.RuneCountInString(truncationMarker)
|
||||
if keep < 0 {
|
||||
keep = 0
|
||||
}
|
||||
reportText = truncateRunes(reportText, keep) + truncationMarker
|
||||
}
|
||||
|
||||
var b strings.Builder
|
||||
b.WriteString(head.String())
|
||||
b.WriteString(fence)
|
||||
b.WriteString("\n")
|
||||
b.WriteString(reportText)
|
||||
b.WriteString("\n")
|
||||
b.WriteString(fence)
|
||||
b.WriteString(footer)
|
||||
|
||||
// Backstop: never exceed the cap even if the fixed sections themselves are
|
||||
// unexpectedly large.
|
||||
return truncateRunes(b.String(), MaxIssueBodyRunes)
|
||||
}
|
||||
|
||||
// metadataTable renders the report's non-empty metadata fields as a Markdown
|
||||
// table with values in inline code (neutralising '|' and @mentions).
|
||||
func metadataTable(rep ingest.Report) string {
|
||||
rows := []struct{ label, value string }{
|
||||
{"App version", rep.AppVersion},
|
||||
{"Platform", rep.Platform},
|
||||
{"OS version", rep.OSVersion},
|
||||
{"Device", rep.Device},
|
||||
{"Client timestamp", rep.ClientTimestamp},
|
||||
}
|
||||
var b strings.Builder
|
||||
b.WriteString("| Field | Value |\n| --- | --- |\n")
|
||||
for _, r := range rows {
|
||||
v := oneLine(r.value)
|
||||
if v == "" {
|
||||
continue
|
||||
}
|
||||
b.WriteString("| ")
|
||||
b.WriteString(r.label)
|
||||
b.WriteString(" | `")
|
||||
b.WriteString(sanitizeCode(v))
|
||||
b.WriteString("` |\n")
|
||||
}
|
||||
return b.String()
|
||||
}
|
||||
|
||||
// fenceFor returns a run of backticks one longer than the longest backtick run in
|
||||
// s (minimum 3), so s can be embedded in a fenced code block without the fence
|
||||
// being closed early by backticks inside the content.
|
||||
func fenceFor(s string) string {
|
||||
longest, run := 0, 0
|
||||
for _, r := range s {
|
||||
if r == '`' {
|
||||
run++
|
||||
if run > longest {
|
||||
longest = run
|
||||
}
|
||||
} else {
|
||||
run = 0
|
||||
}
|
||||
}
|
||||
n := longest + 1
|
||||
if n < 3 {
|
||||
n = 3
|
||||
}
|
||||
return strings.Repeat("`", n)
|
||||
}
|
||||
|
||||
// oneLine collapses a metadata value to a single trimmed line (newlines and other
|
||||
// control characters become spaces), keeping the issue tidy.
|
||||
func oneLine(s string) string {
|
||||
s = strings.Map(func(r rune) rune {
|
||||
if r == '\n' || r == '\r' || r == '\t' {
|
||||
return ' '
|
||||
}
|
||||
if r < 0x20 {
|
||||
return -1 // drop other control characters
|
||||
}
|
||||
return r
|
||||
}, s)
|
||||
return strings.TrimSpace(s)
|
||||
}
|
||||
|
||||
// sanitizeCode makes a value safe to place inside an inline-code span within a
|
||||
// Markdown table cell: backticks (which would close the span) become quotes and
|
||||
// pipes (which would split the cell) become slashes.
|
||||
func sanitizeCode(s string) string {
|
||||
s = strings.ReplaceAll(s, "`", "'")
|
||||
s = strings.ReplaceAll(s, "|", "/")
|
||||
return s
|
||||
}
|
||||
|
||||
// truncateRunes returns s limited to at most n runes (n < 0 is treated as 0).
|
||||
func truncateRunes(s string, n int) string {
|
||||
if n <= 0 {
|
||||
return ""
|
||||
}
|
||||
if utf8.RuneCountInString(s) <= n {
|
||||
return s
|
||||
}
|
||||
i, count := 0, 0
|
||||
for i = range s {
|
||||
if count == n {
|
||||
return s[:i]
|
||||
}
|
||||
count++
|
||||
}
|
||||
return s
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
package publish
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"testing"
|
||||
"unicode/utf8"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/ingest"
|
||||
)
|
||||
|
||||
// reportJSON marshals a Report to the scrubbed-JSON plaintext shape that
|
||||
// crypto.Open yields (what formatIssue receives).
|
||||
func reportJSON(t *testing.T, r ingest.Report) []byte {
|
||||
t.Helper()
|
||||
b, err := json.Marshal(r)
|
||||
if err != nil {
|
||||
t.Fatalf("marshal report: %v", err)
|
||||
}
|
||||
return b
|
||||
}
|
||||
|
||||
// TestFormatIssueWellFormed checks the happy path renders a title and a body that
|
||||
// carries the id, the metadata, and the report text inside a code fence.
|
||||
func TestFormatIssueWellFormed(t *testing.T) {
|
||||
id := "20260703T120000-aabbccddeeff00112233"
|
||||
pt := reportJSON(t, ingest.Report{
|
||||
AppVersion: "1.4.2 (142)",
|
||||
Platform: "android",
|
||||
OSVersion: "Android 14",
|
||||
Device: "Pixel 7",
|
||||
ClientTimestamp: "2026-07-02T12:34:56Z",
|
||||
Report: "NullPointerException in SyncService\nat line 42",
|
||||
})
|
||||
|
||||
title, body := formatIssue(id, pt)
|
||||
|
||||
if !strings.Contains(title, "1.4.2 (142)") || !strings.Contains(title, "android") {
|
||||
t.Errorf("title missing app/platform: %q", title)
|
||||
}
|
||||
for _, want := range []string{
|
||||
id, // report id for traceability
|
||||
"1.4.2 (142)", // app version row
|
||||
"Android 14", // os version row
|
||||
"Pixel 7", // device row
|
||||
"2026-07-02T12:34:56Z", // client timestamp row
|
||||
"NullPointerException in SyncService", // the report text
|
||||
"docs/privacy.md", // scrubbing disclosure
|
||||
"bug-report", // labels footer
|
||||
} {
|
||||
if !strings.Contains(body, want) {
|
||||
t.Errorf("body missing %q\n---\n%s", want, body)
|
||||
}
|
||||
}
|
||||
if utf8.RuneCountInString(body) > MaxIssueBodyRunes {
|
||||
t.Errorf("body length %d exceeds cap %d", utf8.RuneCountInString(body), MaxIssueBodyRunes)
|
||||
}
|
||||
// The report text must sit inside a fenced code block (injection-safe).
|
||||
if !strings.Contains(body, "```\nNullPointerException") {
|
||||
t.Errorf("report text not wrapped in a code fence:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFormatIssueTruncatesOversizeBody is the ADR #6 §2.1 requirement: a report
|
||||
// whose rendered body would exceed 65,536 characters is truncated to fit.
|
||||
func TestFormatIssueTruncatesOversizeBody(t *testing.T) {
|
||||
huge := strings.Repeat("A", MaxIssueBodyRunes*2) // way over the cap
|
||||
pt := reportJSON(t, ingest.Report{
|
||||
AppVersion: "9.9.9",
|
||||
Platform: "android",
|
||||
Report: huge,
|
||||
})
|
||||
|
||||
_, body := formatIssue("id-huge", pt)
|
||||
|
||||
if got := utf8.RuneCountInString(body); got > MaxIssueBodyRunes {
|
||||
t.Fatalf("truncated body length %d exceeds cap %d", got, MaxIssueBodyRunes)
|
||||
}
|
||||
if !strings.Contains(body, "truncated") {
|
||||
t.Errorf("oversize body lacks a truncation marker:\n%s…", body[:200])
|
||||
}
|
||||
// It should still be near the cap (we kept as much as fits), not tiny.
|
||||
if got := utf8.RuneCountInString(body); got < MaxIssueBodyRunes-1024 {
|
||||
t.Errorf("truncated body length %d is far below the cap; too aggressive", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFormatIssueMarkdownInjectionSafe proves backtick runs in the report cannot
|
||||
// break out of the code fence and @mentions are neutralised (kept inside code).
|
||||
func TestFormatIssueMarkdownInjectionSafe(t *testing.T) {
|
||||
// A report containing a triple-backtick run and an @mention.
|
||||
body := "before ``` middle ```` end @maintainer please look"
|
||||
pt := reportJSON(t, ingest.Report{AppVersion: "1", Platform: "android", Report: body})
|
||||
|
||||
_, out := formatIssue("id-x", pt)
|
||||
|
||||
// The fence chosen must be longer than the longest backtick run (4) → ≥5.
|
||||
if !strings.Contains(out, "`````") {
|
||||
t.Errorf("expected a fence of ≥5 backticks to survive the content, got:\n%s", out)
|
||||
}
|
||||
// The @mention is present but only ever inside the code block (we can at least
|
||||
// assert it is not on a line by itself outside a fence — a coarse check: the
|
||||
// content line is indented/enclosed, i.e. the raw text is retained verbatim).
|
||||
if !strings.Contains(out, "@maintainer") {
|
||||
t.Errorf("report content was dropped")
|
||||
}
|
||||
if utf8.RuneCountInString(out) > MaxIssueBodyRunes {
|
||||
t.Errorf("body exceeds cap")
|
||||
}
|
||||
}
|
||||
|
||||
// TestFormatIssueUnparseablePayload: if the decrypted bytes are not the expected
|
||||
// JSON shape, the content is shown verbatim rather than dropped.
|
||||
func TestFormatIssueUnparseablePayload(t *testing.T) {
|
||||
pt := []byte("this is not json, just raw text with a secret [REDACTED_EMAIL]")
|
||||
|
||||
title, body := formatIssue("id-raw", pt)
|
||||
|
||||
if title == "" {
|
||||
t.Error("empty title for unparseable payload")
|
||||
}
|
||||
if !strings.Contains(body, "not json") {
|
||||
t.Errorf("verbatim content missing:\n%s", body)
|
||||
}
|
||||
if !strings.Contains(body, "did not match the expected report schema") {
|
||||
t.Errorf("expected a schema-mismatch note:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
// TestFormatIssuePipeAndBacktickInMetadata: metadata values with '|' or '`' must
|
||||
// not break the Markdown table or the inline-code span.
|
||||
func TestFormatIssueMetadataSanitised(t *testing.T) {
|
||||
pt := reportJSON(t, ingest.Report{
|
||||
AppVersion: "1.0 | weird `build`",
|
||||
Platform: "android",
|
||||
Report: "x",
|
||||
})
|
||||
_, body := formatIssue("id", pt)
|
||||
// The raw '|' and '`' must have been replaced in the rendered cell.
|
||||
if strings.Contains(body, "1.0 | weird `build`") {
|
||||
t.Errorf("metadata not sanitised for table/inline-code:\n%s", body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestFenceFor(t *testing.T) {
|
||||
cases := map[string]int{
|
||||
"no backticks": 3,
|
||||
"one ` here": 3, // longest run 1 → max(3, 2)=3
|
||||
"two `` here": 3, // longest run 2 → max(3,3)=3
|
||||
"three ``` here": 4,
|
||||
"four ```` here": 5,
|
||||
}
|
||||
for in, wantLen := range cases {
|
||||
if got := fenceFor(in); len(got) != wantLen {
|
||||
t.Errorf("fenceFor(%q) len = %d, want %d", in, len(got), wantLen)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestTruncateRunes(t *testing.T) {
|
||||
if got := truncateRunes("héllo", 3); got != "hél" {
|
||||
t.Errorf("truncateRunes multibyte = %q, want %q", got, "hél")
|
||||
}
|
||||
if got := truncateRunes("abc", 10); got != "abc" {
|
||||
t.Errorf("truncateRunes under limit = %q, want abc", got)
|
||||
}
|
||||
if got := truncateRunes("abc", 0); got != "" {
|
||||
t.Errorf("truncateRunes zero = %q, want empty", got)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,493 @@
|
||||
// Package publish turns pending, encrypted bug-reports into GitHub issues on the
|
||||
// LibreMail app repo — the "publish" half of the weekly job (#14). It is handed a
|
||||
// batch of pending report ids by internal/schedule (the Friday-17:00-Central
|
||||
// gate, #13) and, for each id, fetches the still-encrypted frame
|
||||
// (lifecycle.GetPending, #10), decrypts it (internal/crypto, #5), formats it into
|
||||
// a well-formed issue body, and creates a labeled issue via the GitHub REST API.
|
||||
//
|
||||
// # Host-testable, Wasm-deployable
|
||||
//
|
||||
// The package carries no build constraints. The GitHub client is built on the
|
||||
// standard net/http (the #26 patch makes net/http's client work under TinyGo's
|
||||
// js/wasm target, routing through the Cloudflare fetch runtime), so the exact
|
||||
// same code is exercised by `go test` on the host against an httptest mock and
|
||||
// runs in the deployed Worker. All wall-clock behaviour (rate-limit spacing,
|
||||
// backoff sleeps, the ratelimit-reset clock) is injected, so the retry policy is
|
||||
// unit-tested deterministically with no real sleeps.
|
||||
//
|
||||
// # GitHub limits (ADR #6, docs/decisions/labels-and-abuse.md)
|
||||
//
|
||||
// The client encodes ADR #6 §3.2 exactly: serial mutations spaced by ≥1s, honour
|
||||
// Retry-After, wait until x-ratelimit-reset when remaining is 0, a ≥60s floor for
|
||||
// secondary-rate-limit 403s, and full-jittered exponential backoff (base 1s, cap
|
||||
// 60s, ≤5 attempts) on 5xx/network errors. Issues carry the three ADR #6 labels
|
||||
// (bug-report, automated, needs-triage), created if missing. The publish job
|
||||
// applies the ADR #6 §3.1 per-run cap (50) and the 65,536-char body cap.
|
||||
package publish
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"math/rand"
|
||||
"net/http"
|
||||
"strconv"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// GitHub API defaults.
|
||||
const (
|
||||
// defaultBaseURL is the public GitHub REST API root. Overridable with
|
||||
// WithBaseURL so tests can point the client at an httptest server.
|
||||
defaultBaseURL = "https://api.github.com"
|
||||
// apiVersion is sent as X-GitHub-Api-Version, pinning the REST API contract.
|
||||
apiVersion = "2022-11-28"
|
||||
// userAgent identifies this job to GitHub (a User-Agent header is required).
|
||||
userAgent = "LibreMail-Bug-Report-Ingest"
|
||||
)
|
||||
|
||||
// Retry/pacing defaults, from ADR #6 §3.2.
|
||||
const (
|
||||
// defaultMinSpacing is the minimum wall-clock gap between mutating requests
|
||||
// (issue/label creation). ADR #6: "space successful creates by ≥1 second".
|
||||
defaultMinSpacing = 1 * time.Second
|
||||
// defaultBaseBackoff and defaultCapBackoff bound the exponential schedule
|
||||
// min(60s, 1s·2^attempt).
|
||||
defaultBaseBackoff = 1 * time.Second
|
||||
defaultCapBackoff = 60 * time.Second
|
||||
// defaultSecondaryMin is the ≥60s floor for a secondary-rate-limit 403 that
|
||||
// arrives without a Retry-After header.
|
||||
defaultSecondaryMin = 60 * time.Second
|
||||
// defaultMaxAttempts is the per-request attempt cap (ADR #6: "max 5 attempts
|
||||
// per issue"): the initial try plus up to four backoff retries.
|
||||
defaultMaxAttempts = 5
|
||||
)
|
||||
|
||||
// Client is a minimal GitHub REST API client scoped to one owner/repo. It creates
|
||||
// labels and issues and applies ADR #6's pacing, backoff, and rate-limit policy.
|
||||
// It is safe for the serial, single-goroutine use the scheduled publish job makes
|
||||
// of it; it does not add its own concurrency.
|
||||
type Client struct {
|
||||
httpClient *http.Client
|
||||
baseURL string
|
||||
token string
|
||||
owner string
|
||||
repo string
|
||||
|
||||
minSpacing time.Duration
|
||||
baseBackoff time.Duration
|
||||
capBackoff time.Duration
|
||||
secondaryMin time.Duration
|
||||
maxAttempts int
|
||||
|
||||
// Injected wall-clock + randomness seams, so backoff/spacing are deterministic
|
||||
// and instant under test. Defaults use real time and math/rand.
|
||||
now func() time.Time
|
||||
sleep func(context.Context, time.Duration) error
|
||||
randf func() float64
|
||||
|
||||
mu sync.Mutex
|
||||
lastMutation time.Time // time of the last mutating request (for spacing)
|
||||
}
|
||||
|
||||
// ClientOption customises a Client. The wall-clock/rand options exist for tests.
|
||||
type ClientOption func(*Client)
|
||||
|
||||
// WithBaseURL overrides the GitHub API root (tests point it at httptest).
|
||||
func WithBaseURL(u string) ClientOption {
|
||||
return func(c *Client) { c.baseURL = strings.TrimRight(u, "/") }
|
||||
}
|
||||
|
||||
// WithHTTPClient overrides the underlying *http.Client.
|
||||
func WithHTTPClient(h *http.Client) ClientOption {
|
||||
return func(c *Client) { c.httpClient = h }
|
||||
}
|
||||
|
||||
// WithClock injects the time source and sleep function (tests use a virtual
|
||||
// clock so backoff/spacing incur no real delay). Both must be set together.
|
||||
func WithClock(now func() time.Time, sleep func(context.Context, time.Duration) error) ClientOption {
|
||||
return func(c *Client) {
|
||||
if now != nil {
|
||||
c.now = now
|
||||
}
|
||||
if sleep != nil {
|
||||
c.sleep = sleep
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithRand injects the [0,1) source used for full-jitter backoff (tests pin it
|
||||
// to make jittered sleeps deterministic).
|
||||
func WithRand(randf func() float64) ClientOption {
|
||||
return func(c *Client) {
|
||||
if randf != nil {
|
||||
c.randf = randf
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithMaxAttempts overrides the per-request attempt cap (default 5).
|
||||
func WithMaxAttempts(n int) ClientOption {
|
||||
return func(c *Client) {
|
||||
if n > 0 {
|
||||
c.maxAttempts = n
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithSpacing overrides the minimum gap between mutating requests (default 1s).
|
||||
func WithSpacing(d time.Duration) ClientOption {
|
||||
return func(c *Client) { c.minSpacing = d }
|
||||
}
|
||||
|
||||
// NewClient returns a Client for github.com/<owner>/<repo> authenticated with
|
||||
// token (a Worker secret; never log it). Defaults encode ADR #6 §3.2; override
|
||||
// via options (mainly in tests).
|
||||
func NewClient(token, owner, repo string, opts ...ClientOption) *Client {
|
||||
src := rand.New(rand.NewSource(time.Now().UnixNano())) // jitter only; not crypto
|
||||
c := &Client{
|
||||
httpClient: &http.Client{},
|
||||
baseURL: defaultBaseURL,
|
||||
token: token,
|
||||
owner: owner,
|
||||
repo: repo,
|
||||
minSpacing: defaultMinSpacing,
|
||||
baseBackoff: defaultBaseBackoff,
|
||||
capBackoff: defaultCapBackoff,
|
||||
secondaryMin: defaultSecondaryMin,
|
||||
maxAttempts: defaultMaxAttempts,
|
||||
now: time.Now,
|
||||
sleep: sleepCtx,
|
||||
randf: src.Float64,
|
||||
}
|
||||
for _, o := range opts {
|
||||
o(c)
|
||||
}
|
||||
return c
|
||||
}
|
||||
|
||||
// sleepCtx sleeps for d, returning early if ctx is cancelled. It is the default
|
||||
// Client sleep; tests swap in a virtual-clock sleep that returns instantly.
|
||||
func sleepCtx(ctx context.Context, d time.Duration) error {
|
||||
if d <= 0 {
|
||||
return ctx.Err()
|
||||
}
|
||||
t := time.NewTimer(d)
|
||||
defer t.Stop()
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-t.C:
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// Label is a GitHub issue label to ensure-exists and apply. Colour is a 6-hex-
|
||||
// digit string without the leading '#'.
|
||||
type Label struct {
|
||||
Name string
|
||||
Color string
|
||||
Description string
|
||||
}
|
||||
|
||||
// DefaultLabels are the three ADR #6 §1 labels every auto-published issue carries.
|
||||
// Names are the binding contract; colours are advisory (the maintainer may
|
||||
// recolour). All three are applied at issue creation.
|
||||
var DefaultLabels = []Label{
|
||||
{Name: "bug-report", Color: "e99695", Description: "Originated from the app's opt-in debug bug-report ingest pipeline."},
|
||||
{Name: "automated", Color: "c5def5", Description: "Created by the weekly publish job, not written by a human."},
|
||||
{Name: "needs-triage", Color: "fbca04", Description: "Not yet reviewed or confirmed by the maintainer."},
|
||||
}
|
||||
|
||||
// LabelNames returns just the label names, for the issue-creation payload.
|
||||
func LabelNames(labels []Label) []string {
|
||||
names := make([]string, len(labels))
|
||||
for i, l := range labels {
|
||||
names[i] = l.Name
|
||||
}
|
||||
return names
|
||||
}
|
||||
|
||||
// CreatedIssue is the slice of GitHub's issue-creation response the job needs.
|
||||
type CreatedIssue struct {
|
||||
Number int `json:"number"`
|
||||
HTMLURL string `json:"html_url"`
|
||||
}
|
||||
|
||||
// EnsureLabels creates each label that does not yet exist, idempotently (ADR #6:
|
||||
// "create-or-ignore"). An already-existing label (HTTP 422) is treated as
|
||||
// success. Errors for individual labels are joined; callers may treat a failure
|
||||
// as best-effort (issue creation still references the label names).
|
||||
func (c *Client) EnsureLabels(ctx context.Context, labels []Label) error {
|
||||
var errs []error
|
||||
for _, l := range labels {
|
||||
if err := c.ensureLabel(ctx, l); err != nil {
|
||||
errs = append(errs, err)
|
||||
}
|
||||
}
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
// labelRequest / issueRequest are the JSON bodies for the two POST endpoints.
|
||||
type labelRequest struct {
|
||||
Name string `json:"name"`
|
||||
Color string `json:"color"`
|
||||
Description string `json:"description,omitempty"`
|
||||
}
|
||||
|
||||
type issueRequest struct {
|
||||
Title string `json:"title"`
|
||||
Body string `json:"body"`
|
||||
Labels []string `json:"labels,omitempty"`
|
||||
}
|
||||
|
||||
func (c *Client) ensureLabel(ctx context.Context, l Label) error {
|
||||
body, err := json.Marshal(labelRequest{Name: l.Name, Color: l.Color, Description: l.Description})
|
||||
if err != nil {
|
||||
return fmt.Errorf("github: marshal label %q: %w", l.Name, err)
|
||||
}
|
||||
status, respBody, err := c.do(ctx, http.MethodPost, c.repoPath("labels"), body)
|
||||
if err != nil {
|
||||
return fmt.Errorf("github: create label %q: %w", l.Name, err)
|
||||
}
|
||||
switch status {
|
||||
case http.StatusCreated:
|
||||
return nil
|
||||
case http.StatusUnprocessableEntity:
|
||||
// 422 on label creation means the name already exists ("already_exists"):
|
||||
// the idempotent create-or-ignore outcome ADR #6 asks for.
|
||||
return nil
|
||||
default:
|
||||
return fmt.Errorf("github: create label %q: unexpected status %d: %s", l.Name, status, snippet(respBody))
|
||||
}
|
||||
}
|
||||
|
||||
// CreateIssue creates one issue and returns it. Per ADR #6 a report is publishable
|
||||
// only on a confirmed 201 Created, so a nil error from CreateIssue is exactly that
|
||||
// confirmation — the caller marks the report published only then.
|
||||
func (c *Client) CreateIssue(ctx context.Context, title, body string, labels []string) (CreatedIssue, error) {
|
||||
reqBody, err := json.Marshal(issueRequest{Title: title, Body: body, Labels: labels})
|
||||
if err != nil {
|
||||
return CreatedIssue{}, fmt.Errorf("github: marshal issue: %w", err)
|
||||
}
|
||||
status, respBody, err := c.do(ctx, http.MethodPost, c.repoPath("issues"), reqBody)
|
||||
if err != nil {
|
||||
return CreatedIssue{}, fmt.Errorf("github: create issue: %w", err)
|
||||
}
|
||||
if status != http.StatusCreated {
|
||||
return CreatedIssue{}, fmt.Errorf("github: create issue: unexpected status %d: %s", status, snippet(respBody))
|
||||
}
|
||||
var issue CreatedIssue
|
||||
if err := json.Unmarshal(respBody, &issue); err != nil {
|
||||
return CreatedIssue{}, fmt.Errorf("github: create issue: decode response: %w", err)
|
||||
}
|
||||
return issue, nil
|
||||
}
|
||||
|
||||
// repoPath builds "/repos/<owner>/<repo>/<sub>".
|
||||
func (c *Client) repoPath(sub string) string {
|
||||
return "/repos/" + c.owner + "/" + c.repo + "/" + sub
|
||||
}
|
||||
|
||||
// do performs a mutating request with ADR #6 pacing + retry. It returns the final
|
||||
// HTTP status and body. A non-nil error means the request could not be completed
|
||||
// within the attempt budget (a transport error, or an exhausted-retry status); a
|
||||
// non-retryable HTTP status (2xx, or a 4xx like 401/404/422) returns that status
|
||||
// with a nil error so the caller can interpret it.
|
||||
func (c *Client) do(ctx context.Context, method, path string, reqBody []byte) (int, []byte, error) {
|
||||
url := c.baseURL + path
|
||||
var (
|
||||
lastStatus int
|
||||
lastBody []byte
|
||||
lastErr error
|
||||
)
|
||||
for attempt := 0; attempt < c.maxAttempts; attempt++ {
|
||||
// Serialize + space mutations (ADR #6): ensure ≥minSpacing since the last
|
||||
// one before issuing this attempt.
|
||||
if err := c.spaceMutations(ctx); err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
|
||||
req, err := http.NewRequestWithContext(ctx, method, url, bytes.NewReader(reqBody))
|
||||
if err != nil {
|
||||
return 0, nil, err // request construction failure is not retryable
|
||||
}
|
||||
c.setHeaders(req)
|
||||
|
||||
resp, err := c.httpClient.Do(req)
|
||||
c.recordMutation()
|
||||
|
||||
var wait time.Duration
|
||||
var retryable bool
|
||||
if err != nil {
|
||||
// Network/transport error: retry with jittered exponential backoff.
|
||||
lastStatus, lastBody, lastErr = 0, nil, err
|
||||
wait, retryable = c.jitteredBackoff(attempt), true
|
||||
} else {
|
||||
body, readErr := io.ReadAll(resp.Body)
|
||||
resp.Body.Close()
|
||||
if readErr != nil {
|
||||
lastStatus, lastBody, lastErr = 0, nil, readErr
|
||||
wait, retryable = c.jitteredBackoff(attempt), true
|
||||
} else {
|
||||
wait, retryable = c.retryDecision(resp, body, attempt)
|
||||
if !retryable {
|
||||
return resp.StatusCode, body, nil
|
||||
}
|
||||
lastStatus, lastBody = resp.StatusCode, body
|
||||
lastErr = fmt.Errorf("status %d: %s", resp.StatusCode, snippet(body))
|
||||
}
|
||||
}
|
||||
|
||||
if attempt == c.maxAttempts-1 {
|
||||
return lastStatus, lastBody, fmt.Errorf("github: %s %s: gave up after %d attempt(s): %w",
|
||||
method, path, c.maxAttempts, lastErr)
|
||||
}
|
||||
if err := c.sleep(ctx, wait); err != nil {
|
||||
return 0, nil, err
|
||||
}
|
||||
}
|
||||
// Unreachable (maxAttempts ≥ 1), but keeps the compiler happy.
|
||||
return lastStatus, lastBody, lastErr
|
||||
}
|
||||
|
||||
// setHeaders applies the standard GitHub REST headers. The token is a secret and
|
||||
// must never be logged.
|
||||
func (c *Client) setHeaders(req *http.Request) {
|
||||
req.Header.Set("Authorization", "Bearer "+c.token)
|
||||
req.Header.Set("Accept", "application/vnd.github+json")
|
||||
req.Header.Set("X-GitHub-Api-Version", apiVersion)
|
||||
req.Header.Set("Content-Type", "application/json")
|
||||
req.Header.Set("User-Agent", userAgent)
|
||||
}
|
||||
|
||||
// retryDecision encodes ADR #6 §3.2's rate-limit/error policy: it reports how long
|
||||
// to wait before the next attempt and whether the response is retryable at all.
|
||||
//
|
||||
// - 5xx: retryable, jittered exponential backoff.
|
||||
// - 403/429 with Retry-After: retryable, honour it exactly.
|
||||
// - 403/429 with x-ratelimit-remaining: 0: retryable, wait until x-ratelimit-reset.
|
||||
// - 429 otherwise: retryable, jittered exponential backoff.
|
||||
// - 403 secondary rate limit without Retry-After: retryable, ≥60s then exponential.
|
||||
// - 403 without any rate-limit signal (a permission error): NOT retryable.
|
||||
// - everything else (2xx/4xx): not retryable (the caller interprets the status).
|
||||
func (c *Client) retryDecision(resp *http.Response, body []byte, attempt int) (time.Duration, bool) {
|
||||
code := resp.StatusCode
|
||||
if code >= 500 {
|
||||
return c.jitteredBackoff(attempt), true
|
||||
}
|
||||
if code != http.StatusForbidden && code != http.StatusTooManyRequests {
|
||||
return 0, false
|
||||
}
|
||||
|
||||
if d, ok := parseRetryAfter(resp.Header.Get("Retry-After")); ok {
|
||||
return d, true // honour Retry-After exactly (overrides computed backoff)
|
||||
}
|
||||
if resp.Header.Get("X-RateLimit-Remaining") == "0" {
|
||||
if d, ok := c.untilReset(resp.Header.Get("X-RateLimit-Reset")); ok {
|
||||
return d, true
|
||||
}
|
||||
return c.jitteredBackoff(attempt), true // remaining=0 but no usable reset
|
||||
}
|
||||
if code == http.StatusTooManyRequests {
|
||||
return c.jitteredBackoff(attempt), true
|
||||
}
|
||||
// From here code == 403 with no Retry-After and remaining != 0.
|
||||
if isSecondaryRateLimit(body) {
|
||||
d := c.jitteredBackoff(attempt)
|
||||
if d < c.secondaryMin {
|
||||
d = c.secondaryMin // ADR #6: ≥60s floor for secondary limit
|
||||
}
|
||||
return d, true
|
||||
}
|
||||
// A plain 403 (bad token, no push access, ...) is permanent — do not retry.
|
||||
return 0, false
|
||||
}
|
||||
|
||||
// spaceMutations blocks until at least minSpacing has elapsed since the previous
|
||||
// mutating request, implementing ADR #6's "≥1s between mutations".
|
||||
func (c *Client) spaceMutations(ctx context.Context) error {
|
||||
c.mu.Lock()
|
||||
last := c.lastMutation
|
||||
c.mu.Unlock()
|
||||
if last.IsZero() {
|
||||
return nil
|
||||
}
|
||||
if elapsed := c.now().Sub(last); elapsed < c.minSpacing {
|
||||
return c.sleep(ctx, c.minSpacing-elapsed)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// recordMutation stamps the time of the just-issued mutating request.
|
||||
func (c *Client) recordMutation() {
|
||||
c.mu.Lock()
|
||||
c.lastMutation = c.now()
|
||||
c.mu.Unlock()
|
||||
}
|
||||
|
||||
// jitteredBackoff returns a full-jitter delay in [0, min(cap, base·2^attempt)]
|
||||
// (ADR #6: base 1s, cap 60s, full jitter).
|
||||
func (c *Client) jitteredBackoff(attempt int) time.Duration {
|
||||
backoff := c.baseBackoff << attempt // base · 2^attempt
|
||||
if backoff <= 0 || backoff > c.capBackoff {
|
||||
backoff = c.capBackoff
|
||||
}
|
||||
return time.Duration(c.randf() * float64(backoff))
|
||||
}
|
||||
|
||||
// untilReset parses an x-ratelimit-reset epoch-seconds header and returns the
|
||||
// (clamped, +1s buffer) wait until that instant per the injected clock.
|
||||
func (c *Client) untilReset(reset string) (time.Duration, bool) {
|
||||
if reset == "" {
|
||||
return 0, false
|
||||
}
|
||||
epoch, err := strconv.ParseInt(strings.TrimSpace(reset), 10, 64)
|
||||
if err != nil {
|
||||
return 0, false
|
||||
}
|
||||
d := time.Unix(epoch, 0).Sub(c.now())
|
||||
if d < 0 {
|
||||
d = 0
|
||||
}
|
||||
return d + time.Second, true // small buffer so we wake just after reset
|
||||
}
|
||||
|
||||
// parseRetryAfter parses a Retry-After header expressed as integer seconds (the
|
||||
// form GitHub sends). Returns false if absent or unparseable.
|
||||
func parseRetryAfter(v string) (time.Duration, bool) {
|
||||
v = strings.TrimSpace(v)
|
||||
if v == "" {
|
||||
return 0, false
|
||||
}
|
||||
secs, err := strconv.Atoi(v)
|
||||
if err != nil || secs < 0 {
|
||||
return 0, false
|
||||
}
|
||||
return time.Duration(secs) * time.Second, true
|
||||
}
|
||||
|
||||
// isSecondaryRateLimit reports whether a 403 body is GitHub's secondary-rate-limit
|
||||
// message (as opposed to a permission error), so only the former is retried.
|
||||
func isSecondaryRateLimit(body []byte) bool {
|
||||
lower := bytes.ToLower(body)
|
||||
return bytes.Contains(lower, []byte("secondary rate limit")) ||
|
||||
bytes.Contains(lower, []byte("abuse detection"))
|
||||
}
|
||||
|
||||
// snippet returns a short, single-line excerpt of a response body for error
|
||||
// messages (bounded so a huge body cannot bloat a log line).
|
||||
func snippet(body []byte) string {
|
||||
const max = 200
|
||||
s := strings.TrimSpace(string(body))
|
||||
s = strings.ReplaceAll(s, "\n", " ")
|
||||
if len(s) > max {
|
||||
return s[:max] + "…"
|
||||
}
|
||||
return s
|
||||
}
|
||||
@@ -0,0 +1,367 @@
|
||||
package publish
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// --- test doubles: a virtual clock and a scriptable GitHub API mock ---
|
||||
|
||||
// fakeClock is a virtual clock whose time only advances when Sleep is called, so
|
||||
// backoff/spacing are exercised with zero real delay and are assertable.
|
||||
type fakeClock struct {
|
||||
mu sync.Mutex
|
||||
t time.Time
|
||||
sleeps []time.Duration
|
||||
}
|
||||
|
||||
func newFakeClock() *fakeClock {
|
||||
return &fakeClock{t: time.Unix(1_700_000_000, 0).UTC()}
|
||||
}
|
||||
|
||||
func (f *fakeClock) now() time.Time {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return f.t
|
||||
}
|
||||
|
||||
func (f *fakeClock) sleep(_ context.Context, d time.Duration) error {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
if d > 0 {
|
||||
f.t = f.t.Add(d)
|
||||
f.sleeps = append(f.sleeps, d)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (f *fakeClock) sleepCount() int {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
return len(f.sleeps)
|
||||
}
|
||||
|
||||
func (f *fakeClock) lastSleep() time.Duration {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
if len(f.sleeps) == 0 {
|
||||
return 0
|
||||
}
|
||||
return f.sleeps[len(f.sleeps)-1]
|
||||
}
|
||||
|
||||
// scriptedResp is one canned HTTP response for the mock.
|
||||
type scriptedResp struct {
|
||||
status int
|
||||
header map[string]string
|
||||
body string
|
||||
}
|
||||
|
||||
// ghServer is a scriptable mock of the GitHub REST API for POST /labels and
|
||||
// POST /issues. issueScript/labelScript are consumed front-to-back; when empty a
|
||||
// default 201 is returned (and, for issues, a synthetic created-issue body).
|
||||
type ghServer struct {
|
||||
mu sync.Mutex
|
||||
labelReqs []labelRequest
|
||||
createdIssues []issueRequest // recorded only on a 201 response
|
||||
issueAttempts int
|
||||
labelAttempts int
|
||||
issueScript []scriptedResp
|
||||
labelScript []scriptedResp
|
||||
nextIssueNum int
|
||||
authHeaders []string
|
||||
uaHeaders []string
|
||||
apiVersions []string
|
||||
}
|
||||
|
||||
func (g *ghServer) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
g.mu.Lock()
|
||||
defer g.mu.Unlock()
|
||||
|
||||
g.authHeaders = append(g.authHeaders, r.Header.Get("Authorization"))
|
||||
g.uaHeaders = append(g.uaHeaders, r.Header.Get("User-Agent"))
|
||||
g.apiVersions = append(g.apiVersions, r.Header.Get("X-GitHub-Api-Version"))
|
||||
|
||||
switch {
|
||||
case strings.HasSuffix(r.URL.Path, "/labels") && r.Method == http.MethodPost:
|
||||
var lr labelRequest
|
||||
_ = json.NewDecoder(r.Body).Decode(&lr)
|
||||
g.labelReqs = append(g.labelReqs, lr)
|
||||
g.labelAttempts++
|
||||
resp := popScript(&g.labelScript, scriptedResp{status: http.StatusCreated, body: `{}`})
|
||||
writeResp(w, resp)
|
||||
|
||||
case strings.HasSuffix(r.URL.Path, "/issues") && r.Method == http.MethodPost:
|
||||
var ir issueRequest
|
||||
_ = json.NewDecoder(r.Body).Decode(&ir)
|
||||
g.issueAttempts++
|
||||
resp := popScript(&g.issueScript, scriptedResp{})
|
||||
if resp.status == 0 { // default: synthesize a success
|
||||
g.nextIssueNum++
|
||||
n := g.nextIssueNum
|
||||
resp = scriptedResp{
|
||||
status: http.StatusCreated,
|
||||
body: fmt.Sprintf(`{"number":%d,"html_url":"https://github.com/o/r/issues/%d"}`, n, n),
|
||||
}
|
||||
}
|
||||
if resp.status == http.StatusCreated {
|
||||
g.createdIssues = append(g.createdIssues, ir)
|
||||
}
|
||||
writeResp(w, resp)
|
||||
|
||||
default:
|
||||
w.WriteHeader(http.StatusNotFound)
|
||||
}
|
||||
}
|
||||
|
||||
func popScript(q *[]scriptedResp, def scriptedResp) scriptedResp {
|
||||
if len(*q) == 0 {
|
||||
return def
|
||||
}
|
||||
r := (*q)[0]
|
||||
*q = (*q)[1:]
|
||||
return r
|
||||
}
|
||||
|
||||
func writeResp(w http.ResponseWriter, r scriptedResp) {
|
||||
for k, v := range r.header {
|
||||
w.Header().Set(k, v)
|
||||
}
|
||||
w.WriteHeader(r.status)
|
||||
_, _ = w.Write([]byte(r.body))
|
||||
}
|
||||
|
||||
// testClient builds a Client aimed at srvURL with a virtual clock and max-jitter
|
||||
// (randf==1) so exponential backoff is the full computed value, deterministically.
|
||||
func testClient(srvURL string, fc *fakeClock, opts ...ClientOption) *Client {
|
||||
base := []ClientOption{
|
||||
WithBaseURL(srvURL),
|
||||
WithClock(fc.now, fc.sleep),
|
||||
WithRand(func() float64 { return 1.0 }),
|
||||
}
|
||||
return NewClient("test-token", "o", "r", append(base, opts...)...)
|
||||
}
|
||||
|
||||
// --- tests ---
|
||||
|
||||
func TestCreateIssueSuccess(t *testing.T) {
|
||||
g := &ghServer{}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, newFakeClock())
|
||||
|
||||
issue, err := c.CreateIssue(context.Background(), "the title", "the body", []string{"bug-report", "automated"})
|
||||
if err != nil {
|
||||
t.Fatalf("CreateIssue: %v", err)
|
||||
}
|
||||
if issue.Number != 1 || !strings.Contains(issue.HTMLURL, "/issues/1") {
|
||||
t.Errorf("unexpected issue: %+v", issue)
|
||||
}
|
||||
if len(g.createdIssues) != 1 {
|
||||
t.Fatalf("server recorded %d issues, want 1", len(g.createdIssues))
|
||||
}
|
||||
got := g.createdIssues[0]
|
||||
if got.Title != "the title" || got.Body != "the body" {
|
||||
t.Errorf("issue payload = %+v", got)
|
||||
}
|
||||
if len(got.Labels) != 2 || got.Labels[0] != "bug-report" {
|
||||
t.Errorf("labels = %v", got.Labels)
|
||||
}
|
||||
// Required GitHub headers.
|
||||
if g.authHeaders[0] != "Bearer test-token" {
|
||||
t.Errorf("auth header = %q", g.authHeaders[0])
|
||||
}
|
||||
if g.uaHeaders[0] == "" {
|
||||
t.Error("missing User-Agent")
|
||||
}
|
||||
if g.apiVersions[0] != apiVersion {
|
||||
t.Errorf("api version = %q, want %q", g.apiVersions[0], apiVersion)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEnsureLabelsCreateOrIgnore(t *testing.T) {
|
||||
g := &ghServer{
|
||||
// First label already exists (422), the rest are created (201, the default).
|
||||
labelScript: []scriptedResp{{status: http.StatusUnprocessableEntity, body: `{"message":"Validation Failed"}`}},
|
||||
}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, newFakeClock())
|
||||
|
||||
if err := c.EnsureLabels(context.Background(), DefaultLabels); err != nil {
|
||||
t.Fatalf("EnsureLabels: %v", err)
|
||||
}
|
||||
if g.labelAttempts != len(DefaultLabels) {
|
||||
t.Fatalf("label POSTs = %d, want %d", g.labelAttempts, len(DefaultLabels))
|
||||
}
|
||||
// Names and colours must match the ADR contract.
|
||||
names := map[string]string{}
|
||||
for _, l := range g.labelReqs {
|
||||
names[l.Name] = l.Color
|
||||
}
|
||||
for _, want := range DefaultLabels {
|
||||
if names[want.Name] != want.Color {
|
||||
t.Errorf("label %q colour = %q, want %q", want.Name, names[want.Name], want.Color)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryOn5xxThenSucceeds(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
g := &ghServer{issueScript: []scriptedResp{{status: http.StatusInternalServerError, body: "boom"}}}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc)
|
||||
|
||||
if _, err := c.CreateIssue(context.Background(), "t", "b", nil); err != nil {
|
||||
t.Fatalf("CreateIssue should recover after a transient 500: %v", err)
|
||||
}
|
||||
if g.issueAttempts != 2 {
|
||||
t.Errorf("issue attempts = %d, want 2 (one 500, one success)", g.issueAttempts)
|
||||
}
|
||||
// One backoff sleep of the full attempt-0 backoff (1s, randf==1).
|
||||
if fc.sleepCount() != 1 || fc.lastSleep() != 1*time.Second {
|
||||
t.Errorf("backoff sleeps = %v (count %d), want a single 1s sleep", fc.sleeps, fc.sleepCount())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRetryAfterHonoredExactly(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
g := &ghServer{issueScript: []scriptedResp{{
|
||||
status: http.StatusTooManyRequests,
|
||||
header: map[string]string{"Retry-After": "3"},
|
||||
body: "slow down",
|
||||
}}}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc)
|
||||
|
||||
if _, err := c.CreateIssue(context.Background(), "t", "b", nil); err != nil {
|
||||
t.Fatalf("CreateIssue: %v", err)
|
||||
}
|
||||
// Retry-After overrides computed backoff: exactly 3s, not jittered.
|
||||
if fc.lastSleep() != 3*time.Second {
|
||||
t.Errorf("waited %v, want exactly 3s (Retry-After honoured)", fc.lastSleep())
|
||||
}
|
||||
}
|
||||
|
||||
func TestRateLimitResetHonored(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
reset := fc.now().Add(30 * time.Second).Unix()
|
||||
g := &ghServer{issueScript: []scriptedResp{{
|
||||
status: http.StatusForbidden,
|
||||
header: map[string]string{
|
||||
"X-RateLimit-Remaining": "0",
|
||||
"X-RateLimit-Reset": fmt.Sprintf("%d", reset),
|
||||
},
|
||||
body: "rate limited",
|
||||
}}}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc)
|
||||
|
||||
if _, err := c.CreateIssue(context.Background(), "t", "b", nil); err != nil {
|
||||
t.Fatalf("CreateIssue: %v", err)
|
||||
}
|
||||
// Waits until reset (30s) plus the 1s wake buffer.
|
||||
if fc.lastSleep() != 31*time.Second {
|
||||
t.Errorf("waited %v, want 31s (until x-ratelimit-reset + 1s)", fc.lastSleep())
|
||||
}
|
||||
}
|
||||
|
||||
func TestSecondaryRateLimitFloor(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
// A secondary-rate-limit 403 with no Retry-After: floor of 60s applies even
|
||||
// though the attempt-0 exponential backoff would be only ~1s.
|
||||
g := &ghServer{issueScript: []scriptedResp{{
|
||||
status: http.StatusForbidden,
|
||||
body: `{"message":"You have exceeded a secondary rate limit"}`,
|
||||
}}}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc)
|
||||
|
||||
if _, err := c.CreateIssue(context.Background(), "t", "b", nil); err != nil {
|
||||
t.Fatalf("CreateIssue: %v", err)
|
||||
}
|
||||
if fc.lastSleep() < 60*time.Second {
|
||||
t.Errorf("secondary-limit wait = %v, want ≥60s", fc.lastSleep())
|
||||
}
|
||||
}
|
||||
|
||||
func TestPermission403NotRetried(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
// A plain 403 with no rate-limit signals is a permission error: surface it
|
||||
// immediately, do not retry.
|
||||
g := &ghServer{issueScript: []scriptedResp{{
|
||||
status: http.StatusForbidden,
|
||||
body: `{"message":"Resource not accessible by personal access token"}`,
|
||||
}}}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc)
|
||||
|
||||
_, err := c.CreateIssue(context.Background(), "t", "b", nil)
|
||||
if err == nil {
|
||||
t.Fatal("expected an error for a permission 403")
|
||||
}
|
||||
if g.issueAttempts != 1 {
|
||||
t.Errorf("issue attempts = %d, want 1 (permission 403 must not retry)", g.issueAttempts)
|
||||
}
|
||||
if fc.sleepCount() != 0 {
|
||||
t.Errorf("slept %d times on a non-retryable 403, want 0", fc.sleepCount())
|
||||
}
|
||||
}
|
||||
|
||||
func TestMaxAttemptsExhausted(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
// Every attempt fails with 500; with maxAttempts=3 the client gives up after 3.
|
||||
g := &ghServer{issueScript: []scriptedResp{
|
||||
{status: http.StatusBadGateway, body: "1"},
|
||||
{status: http.StatusBadGateway, body: "2"},
|
||||
{status: http.StatusBadGateway, body: "3"},
|
||||
{status: http.StatusBadGateway, body: "4"},
|
||||
}}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc, WithMaxAttempts(3))
|
||||
|
||||
_, err := c.CreateIssue(context.Background(), "t", "b", nil)
|
||||
if err == nil {
|
||||
t.Fatal("expected an error after exhausting attempts")
|
||||
}
|
||||
if g.issueAttempts != 3 {
|
||||
t.Errorf("issue attempts = %d, want 3", g.issueAttempts)
|
||||
}
|
||||
// Two backoffs between three attempts.
|
||||
if fc.sleepCount() != 2 {
|
||||
t.Errorf("backoff sleeps = %d, want 2", fc.sleepCount())
|
||||
}
|
||||
}
|
||||
|
||||
func TestSpacingBetweenMutations(t *testing.T) {
|
||||
fc := newFakeClock()
|
||||
g := &ghServer{}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
c := testClient(srv.URL, fc)
|
||||
|
||||
// Two consecutive creates with no time passing between them: the second must
|
||||
// be spaced by ≥1s (ADR #6).
|
||||
if _, err := c.CreateIssue(context.Background(), "a", "b", nil); err != nil {
|
||||
t.Fatalf("first create: %v", err)
|
||||
}
|
||||
if _, err := c.CreateIssue(context.Background(), "c", "d", nil); err != nil {
|
||||
t.Fatalf("second create: %v", err)
|
||||
}
|
||||
if fc.sleepCount() != 1 || fc.lastSleep() != 1*time.Second {
|
||||
t.Errorf("spacing sleeps = %v, want a single 1s spacing before the 2nd create", fc.sleeps)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,222 @@
|
||||
package publish
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"strings"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/crypto"
|
||||
)
|
||||
|
||||
// DefaultMaxPerRun is ADR #6 §3.1's per-run issue cap: at most 50 issues per
|
||||
// weekly run. If more reports are pending, the oldest 50 are published and the
|
||||
// rest stay pending for the next run (schedule/lifecycle list oldest-first).
|
||||
const DefaultMaxPerRun = 50
|
||||
|
||||
// PendingGetter fetches a pending report's still-encrypted frame by id. It is the
|
||||
// read half of lifecycle.Manager (#10) that the publish job needs; a narrow
|
||||
// interface keeps this package host-testable and free of the Wasm-only storage
|
||||
// backends. *lifecycle.Manager satisfies it via GetPending.
|
||||
type PendingGetter interface {
|
||||
GetPending(ctx context.Context, id string) ([]byte, error)
|
||||
}
|
||||
|
||||
// issueCreator is the slice of the GitHub client the Publisher depends on. *Client
|
||||
// satisfies it; tests use a mock to exercise orchestration (cap, isolation, the
|
||||
// onPublished seam) without an httptest server.
|
||||
type issueCreator interface {
|
||||
EnsureLabels(ctx context.Context, labels []Label) error
|
||||
CreateIssue(ctx context.Context, title, body string, labels []string) (CreatedIssue, error)
|
||||
}
|
||||
|
||||
// Publisher implements schedule.Publisher (#13): given a batch of pending report
|
||||
// ids it decrypts, formats, and creates one labeled GitHub issue per report,
|
||||
// isolating per-report failures. It is the real publisher that replaces
|
||||
// schedule.LogPublisher in the Worker.
|
||||
type Publisher struct {
|
||||
gh issueCreator
|
||||
keyring *crypto.Keyring
|
||||
getter PendingGetter
|
||||
labels []Label
|
||||
maxPerRun int
|
||||
onPublished func(ctx context.Context, id string) error
|
||||
logf func(format string, args ...any)
|
||||
}
|
||||
|
||||
// Option customises a Publisher.
|
||||
type Option func(*Publisher)
|
||||
|
||||
// WithOnPublished sets the post-publish hook. It is invoked once per report,
|
||||
// immediately after that report's issue is confirmed created (a 201), with the
|
||||
// report id. This is the seam #15 wires to lifecycle.MarkPublished to complete
|
||||
// cross-run de-duplication (see the Publish doc). The default is a no-op.
|
||||
func WithOnPublished(fn func(ctx context.Context, id string) error) Option {
|
||||
return func(p *Publisher) {
|
||||
if fn != nil {
|
||||
p.onPublished = fn
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// WithMaxPerRun overrides the per-run issue cap (default DefaultMaxPerRun). A
|
||||
// value ≤ 0 disables the cap.
|
||||
func WithMaxPerRun(n int) Option {
|
||||
return func(p *Publisher) { p.maxPerRun = n }
|
||||
}
|
||||
|
||||
// WithLabels overrides the labels applied to every issue (default DefaultLabels).
|
||||
func WithLabels(labels []Label) Option {
|
||||
return func(p *Publisher) { p.labels = labels }
|
||||
}
|
||||
|
||||
// WithLogger overrides the logger (default log.Printf). Ids are opaque handles
|
||||
// and safe to log; report plaintext and the GitHub token must never be logged.
|
||||
func WithLogger(fn func(format string, args ...any)) Option {
|
||||
return func(p *Publisher) {
|
||||
if fn != nil {
|
||||
p.logf = fn
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// New returns a Publisher that decrypts with kr, reads pending frames from getter,
|
||||
// and creates issues via gh. kr and getter must be non-nil.
|
||||
func New(gh *Client, kr *crypto.Keyring, getter PendingGetter, opts ...Option) *Publisher {
|
||||
return newPublisher(gh, kr, getter, opts...)
|
||||
}
|
||||
|
||||
// newPublisher is New over the issueCreator seam, so tests can inject a mock
|
||||
// client. Exported New takes the concrete *Client the Worker constructs.
|
||||
func newPublisher(gh issueCreator, kr *crypto.Keyring, getter PendingGetter, opts ...Option) *Publisher {
|
||||
p := &Publisher{
|
||||
gh: gh,
|
||||
keyring: kr,
|
||||
getter: getter,
|
||||
labels: DefaultLabels,
|
||||
maxPerRun: DefaultMaxPerRun,
|
||||
onPublished: func(context.Context, string) error { return nil },
|
||||
logf: log.Printf,
|
||||
}
|
||||
for _, o := range opts {
|
||||
o(p)
|
||||
}
|
||||
return p
|
||||
}
|
||||
|
||||
// Publish is the schedule.Publisher entry point. For each pending id (oldest
|
||||
// first, capped at maxPerRun) it runs GetPending → crypto.Open → format →
|
||||
// CreateIssue, applying the ADR #6 labels, then calls the onPublished hook.
|
||||
//
|
||||
// # De-duplication seam (#15)
|
||||
//
|
||||
// This ticket (#14) creates the issues; it deliberately does NOT mark reports
|
||||
// published. Cross-run de-dup is completed by #15, which injects onPublished →
|
||||
// lifecycle.MarkPublished so a published report leaves the pending set and is not
|
||||
// re-published next run. The contract that makes this safe (ADR #6 §3.2,
|
||||
// "mark published only on a confirmed 201") is upheld here: onPublished runs only
|
||||
// after CreateIssue returns success, so #15 only ever marks confirmed-created
|
||||
// reports. Until #15 lands the default onPublished is a no-op, so reports remain
|
||||
// pending and would republish — acceptable pre-#15 and intentional.
|
||||
//
|
||||
// # Failure isolation
|
||||
//
|
||||
// A per-report failure (fetch, decrypt, create-after-retries, or onPublished) is
|
||||
// logged and collected but never aborts the batch: the remaining reports are
|
||||
// still attempted. Publish returns the joined per-report errors (nil if all
|
||||
// succeeded), which the Worker surfaces so the scheduled invocation is recorded
|
||||
// as failed and the maintainer is alerted. A failed report is not marked
|
||||
// published (onPublished not called), so it is retried on a later run.
|
||||
func (p *Publisher) Publish(ctx context.Context, ids []string) error {
|
||||
if len(ids) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
batch := ids
|
||||
capHit := false
|
||||
if p.maxPerRun > 0 && len(batch) > p.maxPerRun {
|
||||
batch = batch[:p.maxPerRun] // oldest-first; remainder drains next run
|
||||
capHit = true
|
||||
}
|
||||
|
||||
// Ensure the three ADR #6 labels exist once per run (idempotent). Best-effort:
|
||||
// a failure is surfaced but does not abort publishing — issue creation still
|
||||
// references the label names, and any missing label is retried next run.
|
||||
var errs []error
|
||||
if err := p.gh.EnsureLabels(ctx, p.labels); err != nil {
|
||||
p.logf("publish: ensure labels (continuing, labels applied best-effort): %v", err)
|
||||
errs = append(errs, fmt.Errorf("ensure labels: %w", err))
|
||||
}
|
||||
|
||||
names := LabelNames(p.labels)
|
||||
published := 0
|
||||
for _, id := range batch {
|
||||
if err := ctx.Err(); err != nil {
|
||||
errs = append(errs, err)
|
||||
break
|
||||
}
|
||||
if err := p.publishOne(ctx, id, names); err != nil {
|
||||
p.logf("publish: report %s failed, left pending: %v", id, err)
|
||||
errs = append(errs, fmt.Errorf("report %s: %w", id, err))
|
||||
continue
|
||||
}
|
||||
published++
|
||||
}
|
||||
|
||||
if capHit {
|
||||
// ADR #6 §3.1: alert the maintainer that the cap was hit and reports were
|
||||
// deferred. The log line is the alert channel until dedicated alerting lands.
|
||||
p.logf("publish: per-run cap %d reached: published %d of %d pending this run; %d deferred to next run",
|
||||
p.maxPerRun, published, len(ids), len(ids)-len(batch))
|
||||
}
|
||||
p.logf("publish: run complete: %d published, %d failed, of %d attempted",
|
||||
published, len(batch)-published, len(batch))
|
||||
|
||||
return errors.Join(errs...)
|
||||
}
|
||||
|
||||
// publishOne runs the full pipeline for a single report id. It returns an error
|
||||
// (never panics) so the caller can isolate this report from the rest of the batch.
|
||||
func (p *Publisher) publishOne(ctx context.Context, id string, labelNames []string) error {
|
||||
frame, err := p.getter.GetPending(ctx, id)
|
||||
if err != nil {
|
||||
return fmt.Errorf("get pending: %w", err)
|
||||
}
|
||||
plaintext, err := crypto.Open(p.keyring, frame)
|
||||
if err != nil {
|
||||
// A decrypt failure is a hard failure for this report (ADR #5): never fall
|
||||
// back to publishing ciphertext. Isolated and surfaced, not fatal.
|
||||
return fmt.Errorf("decrypt: %w", err)
|
||||
}
|
||||
|
||||
title, body := formatIssue(id, plaintext)
|
||||
issue, err := p.gh.CreateIssue(ctx, title, body, labelNames)
|
||||
if err != nil {
|
||||
return fmt.Errorf("create issue: %w", err)
|
||||
}
|
||||
|
||||
// Confirmed 201 Created: invoke the mark-published seam (#15). If it fails the
|
||||
// issue still exists; #15 owns reconciling that (its published marker prevents
|
||||
// a duplicate next run), so we surface it as this report's error.
|
||||
if err := p.onPublished(ctx, id); err != nil {
|
||||
return fmt.Errorf("on-published hook (issue %s created): %w", issue.HTMLURL, err)
|
||||
}
|
||||
p.logf("publish: report %s -> %s", id, issue.HTMLURL)
|
||||
return nil
|
||||
}
|
||||
|
||||
// ParseRepo splits an "owner/repo" target string into its parts. It backs the
|
||||
// configurable target repo (the GITHUB_REPO Worker var); a malformed value is
|
||||
// rejected with an error rather than silently targeting the wrong repo.
|
||||
func ParseRepo(s string) (owner, repo string, err error) {
|
||||
owner, repo, ok := strings.Cut(strings.TrimSpace(s), "/")
|
||||
if !ok || owner == "" || repo == "" || strings.Contains(repo, "/") {
|
||||
return "", "", fmt.Errorf("publish: invalid target repo %q, want \"owner/repo\"", s)
|
||||
}
|
||||
return owner, repo, nil
|
||||
}
|
||||
|
||||
// Compile-time check that the concrete client satisfies the internal seam the
|
||||
// Publisher depends on.
|
||||
var _ issueCreator = (*Client)(nil)
|
||||
@@ -0,0 +1,373 @@
|
||||
package publish
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net/http/httptest"
|
||||
"slices"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/crypto"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/ingest"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/lifecycle"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/schedule"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/storage"
|
||||
)
|
||||
|
||||
// The real Publisher must satisfy the schedule.Publisher seam it replaces.
|
||||
var _ schedule.Publisher = (*Publisher)(nil)
|
||||
|
||||
// --- helpers ---
|
||||
|
||||
func mustKeyring(t *testing.T) *crypto.Keyring {
|
||||
t.Helper()
|
||||
key, err := crypto.GenerateKey()
|
||||
if err != nil {
|
||||
t.Fatalf("GenerateKey: %v", err)
|
||||
}
|
||||
kr, err := crypto.NewKeyring(1, map[uint16][]byte{1: key})
|
||||
if err != nil {
|
||||
t.Fatalf("NewKeyring: %v", err)
|
||||
}
|
||||
return kr
|
||||
}
|
||||
|
||||
func seal(t *testing.T, kr *crypto.Keyring, plaintext []byte) []byte {
|
||||
t.Helper()
|
||||
frame, err := crypto.Seal(kr, plaintext)
|
||||
if err != nil {
|
||||
t.Fatalf("Seal: %v", err)
|
||||
}
|
||||
return frame
|
||||
}
|
||||
|
||||
// fakeGetter is an in-memory PendingGetter (the GetPending seam) keyed by id.
|
||||
type fakeGetter struct {
|
||||
frames map[string][]byte
|
||||
errs map[string]error
|
||||
}
|
||||
|
||||
func (f *fakeGetter) GetPending(_ context.Context, id string) ([]byte, error) {
|
||||
if e := f.errs[id]; e != nil {
|
||||
return nil, e
|
||||
}
|
||||
frame, ok := f.frames[id]
|
||||
if !ok {
|
||||
return nil, errors.New("fakeGetter: no frame for " + id)
|
||||
}
|
||||
return frame, nil
|
||||
}
|
||||
|
||||
// mockCreator is an issueCreator that records created issues and can be scripted
|
||||
// to fail, without any HTTP. It exercises Publisher orchestration in isolation.
|
||||
type mockCreator struct {
|
||||
mu sync.Mutex
|
||||
ensureErr error
|
||||
ensureCalls int
|
||||
created []createRec
|
||||
createFn func(title, body string, labels []string) (CreatedIssue, error)
|
||||
}
|
||||
|
||||
type createRec struct {
|
||||
title, body string
|
||||
labels []string
|
||||
}
|
||||
|
||||
func (m *mockCreator) EnsureLabels(_ context.Context, _ []Label) error {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.ensureCalls++
|
||||
return m.ensureErr
|
||||
}
|
||||
|
||||
func (m *mockCreator) CreateIssue(_ context.Context, title, body string, labels []string) (CreatedIssue, error) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
if m.createFn != nil {
|
||||
iss, err := m.createFn(title, body, labels)
|
||||
if err == nil {
|
||||
m.created = append(m.created, createRec{title, body, labels})
|
||||
}
|
||||
return iss, err
|
||||
}
|
||||
n := len(m.created) + 1
|
||||
m.created = append(m.created, createRec{title, body, labels})
|
||||
return CreatedIssue{Number: n, HTMLURL: fmt.Sprintf("https://github.com/o/r/issues/%d", n)}, nil
|
||||
}
|
||||
|
||||
func recorder() (func(context.Context, string) error, func() []string) {
|
||||
var mu sync.Mutex
|
||||
var ids []string
|
||||
fn := func(_ context.Context, id string) error {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
ids = append(ids, id)
|
||||
return nil
|
||||
}
|
||||
get := func() []string {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
return append([]string(nil), ids...)
|
||||
}
|
||||
return fn, get
|
||||
}
|
||||
|
||||
func sampleReport(id string) ingest.Report {
|
||||
return ingest.Report{
|
||||
AppVersion: "1.4.2 (142)",
|
||||
Platform: "android",
|
||||
OSVersion: "Android 14",
|
||||
Device: "Pixel 7",
|
||||
Report: "crash for " + id,
|
||||
}
|
||||
}
|
||||
|
||||
// --- tests ---
|
||||
|
||||
// TestPublishOneLabeledIssuePerReport is the ticket's core acceptance: N pending
|
||||
// encrypted reports + a test keyring, run through the REAL GitHub client (against
|
||||
// an httptest mock), produce exactly N well-formed, labeled issues and call
|
||||
// onPublished once per success.
|
||||
func TestPublishOneLabeledIssuePerReport(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
ids := []string{"id-1", "id-2", "id-3"}
|
||||
frames := map[string][]byte{}
|
||||
for _, id := range ids {
|
||||
frames[id] = seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
}
|
||||
|
||||
g := &ghServer{}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
client := testClient(srv.URL, newFakeClock())
|
||||
|
||||
onPub, published := recorder()
|
||||
pub := New(client, kr, &fakeGetter{frames: frames}, WithOnPublished(onPub))
|
||||
|
||||
if err := pub.Publish(context.Background(), ids); err != nil {
|
||||
t.Fatalf("Publish: %v", err)
|
||||
}
|
||||
|
||||
if len(g.createdIssues) != len(ids) {
|
||||
t.Fatalf("created %d issues, want %d", len(g.createdIssues), len(ids))
|
||||
}
|
||||
wantLabels := LabelNames(DefaultLabels)
|
||||
for i, iss := range g.createdIssues {
|
||||
if !slices.Equal(iss.Labels, wantLabels) {
|
||||
t.Errorf("issue %d labels = %v, want %v", i, iss.Labels, wantLabels)
|
||||
}
|
||||
if iss.Title == "" || iss.Body == "" {
|
||||
t.Errorf("issue %d not well-formed: %+v", i, iss)
|
||||
}
|
||||
}
|
||||
// The three labels are ensured once per run (idempotent create-or-ignore).
|
||||
if g.labelAttempts != len(DefaultLabels) {
|
||||
t.Errorf("label POSTs = %d, want %d", g.labelAttempts, len(DefaultLabels))
|
||||
}
|
||||
// onPublished called for each success, in order.
|
||||
if got := published(); !slices.Equal(got, ids) {
|
||||
t.Errorf("onPublished ids = %v, want %v", got, ids)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishRespectsPerRunCap: given more than the cap, exactly cap issues are
|
||||
// created (oldest-first) and the rest are left for the next run (ADR #6 §3.1).
|
||||
func TestPublishRespectsPerRunCap(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
const n = DefaultMaxPerRun + 5
|
||||
ids := make([]string, n)
|
||||
frames := map[string][]byte{}
|
||||
for i := range ids {
|
||||
id := fmt.Sprintf("id-%03d", i)
|
||||
ids[i] = id
|
||||
frames[id] = seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
}
|
||||
|
||||
mc := &mockCreator{}
|
||||
onPub, published := recorder()
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames}, WithOnPublished(onPub))
|
||||
|
||||
if err := pub.Publish(context.Background(), ids); err != nil {
|
||||
t.Fatalf("Publish: %v", err)
|
||||
}
|
||||
if len(mc.created) != DefaultMaxPerRun {
|
||||
t.Errorf("created %d issues, want the cap %d", len(mc.created), DefaultMaxPerRun)
|
||||
}
|
||||
if got := len(published()); got != DefaultMaxPerRun {
|
||||
t.Errorf("onPublished called %d times, want the cap %d", got, DefaultMaxPerRun)
|
||||
}
|
||||
// The published ids are the oldest (first) DefaultMaxPerRun.
|
||||
if got := published(); !slices.Equal(got, ids[:DefaultMaxPerRun]) {
|
||||
t.Errorf("published set is not the oldest %d ids", DefaultMaxPerRun)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishIsolatesDecryptFailure: a report that fails to decrypt is isolated
|
||||
// and surfaced, the others still publish, and onPublished is NOT called for it.
|
||||
func TestPublishIsolatesDecryptFailure(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
wrongKr := mustKeyring(t) // different key material -> ErrAuth on Open
|
||||
|
||||
ids := []string{"good-1", "bad", "good-2"}
|
||||
frames := map[string][]byte{
|
||||
"good-1": seal(t, kr, reportJSON(t, sampleReport("good-1"))),
|
||||
"bad": seal(t, wrongKr, reportJSON(t, sampleReport("bad"))), // sealed under the wrong key
|
||||
"good-2": seal(t, kr, reportJSON(t, sampleReport("good-2"))),
|
||||
}
|
||||
|
||||
mc := &mockCreator{}
|
||||
onPub, published := recorder()
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames}, WithOnPublished(onPub))
|
||||
|
||||
err := pub.Publish(context.Background(), ids)
|
||||
if err == nil {
|
||||
t.Fatal("expected a surfaced (non-nil) error for the decrypt failure")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "bad") || !strings.Contains(err.Error(), "decrypt") {
|
||||
t.Errorf("error should name the failed report and reason: %v", err)
|
||||
}
|
||||
if len(mc.created) != 2 {
|
||||
t.Errorf("created %d issues, want 2 (the decrypt failure isolated)", len(mc.created))
|
||||
}
|
||||
if got := published(); !slices.Equal(got, []string{"good-1", "good-2"}) {
|
||||
t.Errorf("onPublished ids = %v, want the two good reports only", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishIsolatesCreateFailure: a persistent create failure on one report
|
||||
// isolates it (others publish) and does NOT call onPublished for it.
|
||||
func TestPublishIsolatesCreateFailure(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
ids := []string{"id-1", "id-2", "id-3"}
|
||||
frames := map[string][]byte{}
|
||||
for _, id := range ids {
|
||||
frames[id] = seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
}
|
||||
|
||||
calls := 0
|
||||
mc := &mockCreator{createFn: func(_, _ string, _ []string) (CreatedIssue, error) {
|
||||
calls++
|
||||
if calls == 2 { // the second report fails permanently (retries exhausted)
|
||||
return CreatedIssue{}, errors.New("github: create issue: gave up after 5 attempts")
|
||||
}
|
||||
return CreatedIssue{Number: calls, HTMLURL: fmt.Sprintf("u/%d", calls)}, nil
|
||||
}}
|
||||
onPub, published := recorder()
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames}, WithOnPublished(onPub))
|
||||
|
||||
err := pub.Publish(context.Background(), ids)
|
||||
if err == nil {
|
||||
t.Fatal("expected a surfaced error for the create failure")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "id-2") {
|
||||
t.Errorf("error should name the failed report id-2: %v", err)
|
||||
}
|
||||
if len(mc.created) != 2 {
|
||||
t.Errorf("created %d issues, want 2 (id-1 and id-3)", len(mc.created))
|
||||
}
|
||||
if got := published(); slices.Contains(got, "id-2") {
|
||||
t.Errorf("onPublished must not be called for the failed report; got %v", got)
|
||||
}
|
||||
if got := published(); !slices.Equal(got, []string{"id-1", "id-3"}) {
|
||||
t.Errorf("onPublished ids = %v, want id-1 and id-3", got)
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishSurfacesOnPublishedError: if the #15 mark-published hook errors, the
|
||||
// issue was still created but the failure is surfaced (so #15 can reconcile) and
|
||||
// the rest of the batch is unaffected.
|
||||
func TestPublishSurfacesOnPublishedError(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
ids := []string{"id-1", "id-2"}
|
||||
frames := map[string][]byte{}
|
||||
for _, id := range ids {
|
||||
frames[id] = seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
}
|
||||
|
||||
mc := &mockCreator{}
|
||||
pub := newPublisher(mc, kr, &fakeGetter{frames: frames},
|
||||
WithOnPublished(func(_ context.Context, id string) error {
|
||||
if id == "id-1" {
|
||||
return errors.New("mark published failed")
|
||||
}
|
||||
return nil
|
||||
}))
|
||||
|
||||
err := pub.Publish(context.Background(), ids)
|
||||
if err == nil {
|
||||
t.Fatal("expected the onPublished error to be surfaced")
|
||||
}
|
||||
if !strings.Contains(err.Error(), "id-1") {
|
||||
t.Errorf("error should name the failing report: %v", err)
|
||||
}
|
||||
// Both issues were created (the hook runs only after a confirmed create).
|
||||
if len(mc.created) != 2 {
|
||||
t.Errorf("created %d issues, want 2 (the hook failure does not undo creation)", len(mc.created))
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublishEmptyBatchNoop: an empty batch does nothing (not even label ensure).
|
||||
func TestPublishEmptyBatchNoop(t *testing.T) {
|
||||
kr := mustKeyring(t)
|
||||
mc := &mockCreator{}
|
||||
pub := newPublisher(mc, kr, &fakeGetter{})
|
||||
if err := pub.Publish(context.Background(), nil); err != nil {
|
||||
t.Fatalf("Publish(nil): %v", err)
|
||||
}
|
||||
if mc.ensureCalls != 0 || len(mc.created) != 0 {
|
||||
t.Errorf("empty batch did work: ensureCalls=%d created=%d", mc.ensureCalls, len(mc.created))
|
||||
}
|
||||
}
|
||||
|
||||
// TestPublisherWorksWithLifecycleManager proves *lifecycle.Manager satisfies the
|
||||
// PendingGetter seam and the whole pipeline runs end to end off the real store.
|
||||
func TestPublisherWorksWithLifecycleManager(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
kr := mustKeyring(t)
|
||||
store := storage.NewMemoryStore()
|
||||
ids := []string{"20260101T000000-aaaa", "20260102T000000-bbbb"}
|
||||
for _, id := range ids {
|
||||
frame := seal(t, kr, reportJSON(t, sampleReport(id)))
|
||||
if err := store.Put(ctx, storage.ReportKey(storage.StatusPending, id), frame); err != nil {
|
||||
t.Fatalf("seed: %v", err)
|
||||
}
|
||||
}
|
||||
manager := lifecycle.New(store)
|
||||
pending, err := manager.ListPending(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("ListPending: %v", err)
|
||||
}
|
||||
|
||||
g := &ghServer{}
|
||||
srv := httptest.NewServer(g)
|
||||
defer srv.Close()
|
||||
client := testClient(srv.URL, newFakeClock())
|
||||
|
||||
onPub, published := recorder()
|
||||
pub := New(client, kr, manager, WithOnPublished(onPub))
|
||||
|
||||
if err := pub.Publish(ctx, pending); err != nil {
|
||||
t.Fatalf("Publish: %v", err)
|
||||
}
|
||||
if len(g.createdIssues) != len(ids) {
|
||||
t.Errorf("created %d issues, want %d", len(g.createdIssues), len(ids))
|
||||
}
|
||||
if got := published(); !slices.Equal(got, pending) {
|
||||
t.Errorf("onPublished ids = %v, want %v", got, pending)
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseRepo(t *testing.T) {
|
||||
owner, repo, err := ParseRepo("JMR-dev/LibreMail")
|
||||
if err != nil || owner != "JMR-dev" || repo != "LibreMail" {
|
||||
t.Errorf("ParseRepo ok case = (%q,%q,%v)", owner, repo, err)
|
||||
}
|
||||
for _, bad := range []string{"", "noslash", "a/b/c", "/LibreMail", "JMR-dev/"} {
|
||||
if _, _, err := ParseRepo(bad); err == nil {
|
||||
t.Errorf("ParseRepo(%q) = nil error, want failure", bad)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -79,8 +79,9 @@ func (s *WorkerSink) loadKeyring() (*crypto.Keyring, error) {
|
||||
}
|
||||
|
||||
// ReadSecret reads a Cloudflare Secrets Store secret by binding name, for use
|
||||
// outside this package (e.g. the Worker admin backend reading AdminTokenBinding,
|
||||
// #11). Like getSecret it must be called within a request handler, and the value
|
||||
// outside this package: the Worker admin backend reads AdminTokenBinding (#11)
|
||||
// and the scheduled publish handler reads GITHUB_TOKEN (#14) through it. Like
|
||||
// getSecret it must be called within a request/scheduled handler, and the value
|
||||
// must never be logged or echoed. It errors if the binding is unbound or empty.
|
||||
func ReadSecret(binding string) ([]byte, error) {
|
||||
return getSecret(binding)
|
||||
|
||||
+88
-24
@@ -9,23 +9,39 @@
|
||||
// (keeping this ticket's diff off the shared fetch wiring).
|
||||
//
|
||||
// Compiled only into the js/wasm Worker; excluded from host builds and tests. All
|
||||
// non-trivial logic — the DST timezone gate and the list→publish orchestration —
|
||||
// lives in internal/schedule, which is build-tag-free and host-tested; this file
|
||||
// is the thin runtime adapter that binds the R2-backed lifecycle Manager to it.
|
||||
// non-trivial logic — the DST timezone gate (internal/schedule) and the
|
||||
// decrypt→format→create publish pipeline (internal/publish) — lives in
|
||||
// build-tag-free, host-tested packages; this file is the thin runtime adapter
|
||||
// that binds the R2-backed lifecycle Manager, the encryption keyring, and the
|
||||
// GitHub token/target-repo (from Cloudflare bindings) to them.
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/syumai/workers/cloudflare"
|
||||
"github.com/syumai/workers/cloudflare/cron"
|
||||
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/crypto"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/lifecycle"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/publish"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/schedule"
|
||||
"github.com/JMR-dev/LibreMail-Bug-Report-Ingest/internal/storage"
|
||||
)
|
||||
|
||||
// Cloudflare binding/var names for the publish step (#14), declared in
|
||||
// wrangler.jsonc. The token is a Secrets Store secret (read via binding.get(),
|
||||
// like the encryption keyring); the target repo is a plain var so it is
|
||||
// configurable without a redeploy of code.
|
||||
const (
|
||||
githubTokenBinding = "GITHUB_TOKEN" // Secrets Store secret holding the PAT
|
||||
githubRepoVar = "GITHUB_REPO" // "owner/repo", e.g. "JMR-dev/LibreMail"
|
||||
defaultTargetRepo = "JMR-dev/LibreMail"
|
||||
)
|
||||
|
||||
// init registers the Cron Trigger handler with the syumai/workers runtime. The
|
||||
// cron package's own init installs the JS "runScheduler" binding; this call just
|
||||
// selects which Task that binding runs.
|
||||
@@ -37,35 +53,83 @@ func init() {
|
||||
//
|
||||
// Cloudflare Cron Triggers are UTC-only, so wrangler.jsonc schedules BOTH Friday
|
||||
// UTC hours that can be 17:00 America/Chicago — 22:00 UTC during CDT (summer) and
|
||||
// 23:00 UTC during CST (winter) — and this handler defers the timezone decision
|
||||
// to schedule.Run. Run gates on the event's scheduled time so only the fire that
|
||||
// is actually 17:00 Central lists the pending reports (via the R2-backed
|
||||
// lifecycle Manager, #10) and hands their ids to the publish step; the sibling
|
||||
// fire is a no-op. Net effect: publishing runs exactly once each Friday, correct
|
||||
// across the DST boundary, from two static UTC crons.
|
||||
// 23:00 UTC during CST (winter) — and the timezone decision is schedule's
|
||||
// IsFriday1700Central gate. Only the fire that is actually 17:00 Central lists
|
||||
// the pending reports (via the R2-backed lifecycle Manager, #10) and publishes
|
||||
// them as GitHub issues; the sibling fire is a no-op. Net effect: publishing runs
|
||||
// exactly once each Friday, correct across the DST boundary, from two static UTC
|
||||
// crons.
|
||||
//
|
||||
// The gate is checked up front so the sibling (no-op) fire does no Secrets Store
|
||||
// I/O: the GitHub token and encryption keyring are read only when the gate is
|
||||
// open. schedule.Run re-applies the same gate authoritatively.
|
||||
func runWeeklyTrigger(ctx context.Context) error {
|
||||
event, err := cron.NewEvent(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Resolve the R2 bucket binding (synchronous, no I/O). It goes unused on the
|
||||
// sibling fire, where Run's gate is closed; that happens at most once a week.
|
||||
store, err := storage.NewR2Store(storage.BucketBinding)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
manager := lifecycle.New(store)
|
||||
|
||||
// schedule.LogPublisher is the no-op default seam; the publish job (#14) will
|
||||
// replace it with the real GitHub-issue publisher without touching this file.
|
||||
ran, err := schedule.Run(ctx, event.ScheduledTime, manager, schedule.LogPublisher{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !ran {
|
||||
// Fast path for the sibling cron: skip the publish wiring (and its secret
|
||||
// reads) entirely. The other Friday fire covers this week.
|
||||
if !schedule.IsFriday1700Central(event.ScheduledTime) {
|
||||
log.Printf("schedule: cron %q fired for %s, not Friday 17:00 America/Chicago; the sibling cron covers this week",
|
||||
event.Cron, event.ScheduledTime.Format(time.RFC3339))
|
||||
return nil
|
||||
}
|
||||
|
||||
manager, publisher, err := buildPublish(ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("schedule: build publisher: %w", err)
|
||||
}
|
||||
|
||||
// schedule.Run re-checks the gate (open here) and hands the pending ids to the
|
||||
// publisher, which decrypts, formats, and creates one labeled GitHub issue per
|
||||
// report. Cross-run de-dup (marking reports published) is #15's job, wired via
|
||||
// the publisher's onPublished hook; #14 leaves that hook at its default no-op.
|
||||
if _, err := schedule.Run(ctx, event.ScheduledTime, manager, publisher); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// buildPublish assembles the pending-report lister and the real publish.Publisher
|
||||
// from the Worker's bindings: the R2-backed lifecycle Manager (source of pending
|
||||
// ids + encrypted frames), the AES-256 keyring and GitHub token from Secrets
|
||||
// Store, and the target repo from a plain var. It performs the (async) secret
|
||||
// reads, so it is only called on the gate-open fire. The Manager is returned
|
||||
// separately so schedule.Run uses the same instance as both lister and getter.
|
||||
func buildPublish(_ context.Context) (*lifecycle.Manager, *publish.Publisher, error) {
|
||||
store, err := storage.NewR2Store(storage.BucketBinding)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("r2 binding: %w", err)
|
||||
}
|
||||
manager := lifecycle.New(store)
|
||||
|
||||
keyringRaw, err := storage.ReadSecret(storage.KeyringBinding)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("read keyring secret: %w", err)
|
||||
}
|
||||
keyring, err := crypto.ParseKeyring(keyringRaw)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("parse keyring: %w", err)
|
||||
}
|
||||
|
||||
tokenRaw, err := storage.ReadSecret(githubTokenBinding)
|
||||
if err != nil {
|
||||
return nil, nil, fmt.Errorf("read github token secret: %w", err)
|
||||
}
|
||||
|
||||
repo := cloudflare.Getenv(githubRepoVar)
|
||||
if repo == "" {
|
||||
repo = defaultTargetRepo
|
||||
}
|
||||
owner, name, err := publish.ParseRepo(repo)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
client := publish.NewClient(string(tokenRaw), owner, name)
|
||||
// onPublished is left at its default no-op: #15 wires it to
|
||||
// lifecycle.MarkPublished to complete cross-run de-duplication.
|
||||
return manager, publish.New(client, keyring, manager), nil
|
||||
}
|
||||
|
||||
+27
-6
@@ -29,12 +29,20 @@
|
||||
}
|
||||
],
|
||||
|
||||
// Cloudflare Secrets Store secret holding the versioned encryption keyring
|
||||
// (ADR #5, Key custody): a single JSON secret {active, keys{ver: base64-32B}}.
|
||||
// "binding" is read at runtime via env.BUGREPORT_ENC_KEYRING.get()
|
||||
// (internal/storage.KeyringBinding). Replace "<store-id>" with the account's
|
||||
// Secrets Store id at deploy time; it is not needed for `pnpm run build`
|
||||
// (Wasm compile) or the devserver, so CI does not require it.
|
||||
// Cloudflare Secrets Store secrets. "binding" is the JS var the Worker reads at
|
||||
// runtime via env.<binding>.get(). Replace "<store-id>" with the account's
|
||||
// Secrets Store id at deploy time; secrets are not needed for `pnpm run build`
|
||||
// (Wasm compile) or the devserver, so CI does not require them.
|
||||
//
|
||||
// - BUGREPORT_ENC_KEYRING (ADR #5, Key custody): the versioned encryption
|
||||
// keyring, a single JSON secret {active, keys{ver: base64-32B}}, read via
|
||||
// internal/storage.KeyringBinding.
|
||||
// - ADMIN_TOKEN (#11): the maintainer admin API shared secret (see the
|
||||
// detailed note on its entry below).
|
||||
// - GITHUB_TOKEN (#14): the token the weekly publish job authenticates to the
|
||||
// GitHub API with to create issues on the target repo. Least privilege: a
|
||||
// fine-grained PAT with Issues: read/write (and Metadata: read) on
|
||||
// JMR-dev/LibreMail. Read via worker/scheduled_wasm.go githubTokenBinding.
|
||||
"secrets_store_secrets": [
|
||||
{
|
||||
"binding": "BUGREPORT_ENC_KEYRING",
|
||||
@@ -53,9 +61,22 @@
|
||||
"binding": "ADMIN_TOKEN",
|
||||
"store_id": "<store-id>",
|
||||
"secret_name": "bugreport-admin-token"
|
||||
},
|
||||
{
|
||||
"binding": "GITHUB_TOKEN",
|
||||
"store_id": "<store-id>",
|
||||
"secret_name": "github-token"
|
||||
}
|
||||
],
|
||||
|
||||
// Plain (non-secret) vars. GITHUB_REPO is the "owner/repo" the weekly publish
|
||||
// job (#14) files issues on; it is configurable here without a code change and
|
||||
// defaults to "JMR-dev/LibreMail" in worker/scheduled_wasm.go if unset. Read via
|
||||
// cloudflare.Getenv(GITHUB_REPO).
|
||||
"vars": {
|
||||
"GITHUB_REPO": "JMR-dev/LibreMail"
|
||||
},
|
||||
|
||||
// Cron Triggers for the weekly publish job (#13): "Friday 17:00 America/Chicago
|
||||
// (Central), DST-correct". Cloudflare evaluates crons in UTC only and has no
|
||||
// timezone support, and 17:00 Central is a different UTC hour depending on DST:
|
||||
|
||||
Reference in New Issue
Block a user