From f9ec7978b81f6911f23442095812b4a5487cca6a Mon Sep 17 00:00:00 2001 From: Jason Ross Date: Thu, 2 Jul 2026 16:25:55 -0500 Subject: [PATCH] Add internal/publish: the real schedule.Publisher that turns pending encrypted reports into labeled GitHub issues. MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - GitHub REST client on net/http (host-testable via httptest; works under TinyGo js/wasm per #26). Encodes ADR #6 §3.2: serial mutations spaced >=1s, honour Retry-After, wait until x-ratelimit-reset, >=60s floor for secondary-limit 403s, full-jitter exponential backoff (base 1s, cap 60s, <=5 attempts). Ensures the three ADR #6 labels (create-or-ignore). - Publisher: GetPending -> crypto.Open -> format -> CreateIssue per id, with the ADR #6 per-run cap (50) and 65,536-char body cap (truncate). Per-report failures are isolated and surfaced, never abort the batch. - onPublished(ctx, id) seam, called only after a confirmed 201, default no-op: #15 wires it to lifecycle.MarkPublished to complete cross-run de-dup. #14 does not implement the mark-published transition. - Issue body wraps report free-text in a length-adaptive code fence and metadata in inline code, neutralising Markdown/@mention injection. - Worker: swap schedule.LogPublisher for the real publisher in scheduled_wasm.go; read GITHUB_TOKEN (Secrets Store) + GITHUB_REPO (var); pre-gate so the sibling cron fire does no secret I/O. worker/main.go untouched. Export storage.GetSecret for the token read. - wrangler.jsonc: add GITHUB_TOKEN secret + GITHUB_REPO var (my keys only). Tests (host, httptest mock, virtual clock): N reports -> N labeled issues + onPublished per success; >65,536-char body truncated; transient 5xx and Retry-After retried per policy; persistent failure isolated (no onPublished); permission 403 not retried; per-run cap; decrypt failure isolated. go vet + go test ./... green; GOOS=js GOARCH=wasm build compiles. Closes #14 Co-Authored-By: Claude Opus 4.8 --- internal/publish/format.go | 204 +++++++++++ internal/publish/format_test.go | 170 +++++++++ internal/publish/github.go | 493 +++++++++++++++++++++++++++ internal/publish/github_test.go | 367 ++++++++++++++++++++ internal/publish/publish.go | 222 ++++++++++++ internal/publish/publish_test.go | 373 ++++++++++++++++++++ internal/storage/worker_sink_wasm.go | 5 +- worker/scheduled_wasm.go | 112 ++++-- wrangler.jsonc | 33 +- 9 files changed, 1947 insertions(+), 32 deletions(-) create mode 100644 internal/publish/format.go create mode 100644 internal/publish/format_test.go create mode 100644 internal/publish/github.go create mode 100644 internal/publish/github_test.go create mode 100644 internal/publish/publish.go create mode 100644 internal/publish/publish_test.go diff --git a/internal/publish/format.go b/internal/publish/format.go new file mode 100644 index 0000000..01eb4be --- /dev/null +++ b/internal/publish/format.go @@ -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 +} diff --git a/internal/publish/format_test.go b/internal/publish/format_test.go new file mode 100644 index 0000000..953616c --- /dev/null +++ b/internal/publish/format_test.go @@ -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) + } +} diff --git a/internal/publish/github.go b/internal/publish/github.go new file mode 100644 index 0000000..d64a456 --- /dev/null +++ b/internal/publish/github.go @@ -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// 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///". +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 +} diff --git a/internal/publish/github_test.go b/internal/publish/github_test.go new file mode 100644 index 0000000..1071f27 --- /dev/null +++ b/internal/publish/github_test.go @@ -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) + } +} diff --git a/internal/publish/publish.go b/internal/publish/publish.go new file mode 100644 index 0000000..a8672a4 --- /dev/null +++ b/internal/publish/publish.go @@ -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) diff --git a/internal/publish/publish_test.go b/internal/publish/publish_test.go new file mode 100644 index 0000000..7610911 --- /dev/null +++ b/internal/publish/publish_test.go @@ -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) + } + } +} diff --git a/internal/storage/worker_sink_wasm.go b/internal/storage/worker_sink_wasm.go index fbccb63..af2b822 100644 --- a/internal/storage/worker_sink_wasm.go +++ b/internal/storage/worker_sink_wasm.go @@ -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) diff --git a/worker/scheduled_wasm.go b/worker/scheduled_wasm.go index 026f31c..c731231 100644 --- a/worker/scheduled_wasm.go +++ b/worker/scheduled_wasm.go @@ -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 +} diff --git a/wrangler.jsonc b/wrangler.jsonc index d740116..639dc3b 100644 --- a/wrangler.jsonc +++ b/wrangler.jsonc @@ -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 "" 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..get(). Replace "" 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": "", "secret_name": "bugreport-admin-token" + }, + { + "binding": "GITHUB_TOKEN", + "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: