diff --git a/docs/configuration.md b/docs/configuration.md index 92c0422..a59a0d8 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -350,6 +350,15 @@ Normalization notes for `v0.3.0`: - `mode = "github_releases"` persists `content_items` as `release` - `mode = "rss"` persists `content_items` as `changelog_entry` +#### Changelog runtime notes for `v0.3.0` +For the first live changelog runtime: + +- `mode = "github_releases"` calls the GitHub releases API for each configured repository and normalizes non-draft releases into `content_items` with `item_type = "release"` +- `mode = "rss"` fetches the configured RSS feed URL and normalizes accepted RSS `` entries into `content_items` with `item_type = "changelog_entry"` +- reruns rely on the shared `content_items` dedupe contract, so previously seen release or feed entries are not inserted again +- release-delta heuristics, Atom-specific handling, and non-GitHub changelog providers are intentionally deferred + + ### Deferred source-family validation The current runtime already reserves `app_store_ios`, `app_store_android`, and `social` as source kinds, but their detailed binding-shape validation is intentionally deferred until those fetchers are implemented. The config contract therefore freezes the source kinds now without overclaiming runtime semantics that do not exist yet. diff --git a/internal/fetch/changelog.go b/internal/fetch/changelog.go new file mode 100644 index 0000000..26ae55f --- /dev/null +++ b/internal/fetch/changelog.go @@ -0,0 +1,727 @@ +package fetch + +import ( + "context" + "encoding/json" + "encoding/xml" + "fmt" + "io" + "net/http" + "net/url" + "os" + "strconv" + "strings" + "time" + + "github.com/CrisSTEM/signalscope/internal/config" + "github.com/CrisSTEM/signalscope/internal/storage" +) + +type changelogBindingMode string + +const ( + SourceKindChangelog = "changelog" + + changelogModeGitHubReleases changelogBindingMode = "github_releases" + changelogModeRSS changelogBindingMode = "rss" + + ContentItemTypeRelease = "release" + ContentItemTypeChangelogEntry = "changelog_entry" + + changelogProviderGitHubReleases = "github_releases" + changelogProviderRSS = "rss" + + defaultChangelogRSSUserAgent = "signalscope-fetch/0.3.0" + + changelogGitHubReleasesPageSize = 100 + changelogGitHubMaxPages = 1000 +) + +type ChangelogFetcherOptions struct { + GitHubBaseURL string + GitHubHTTPClient *http.Client + GitHubUserAgent string + GitHubToken string + + RSSHTTPClient *http.Client + RSSUserAgent string + + Now func() time.Time +} + +type ChangelogFetcher struct { + githubBaseURL *url.URL + githubHTTPClient *http.Client + githubUserAgent string + githubToken string + + rssHTTPClient *http.Client + rssUserAgent string + + now func() time.Time +} + +type changelogBindingScope struct { + Mode changelogBindingMode + Repositories []githubRepositoryRef + FeedURL *url.URL +} + +type changelogGitHubRelease struct { + ID int64 `json:"id"` + TagName string `json:"tag_name"` + Name string `json:"name"` + Body string `json:"body"` + HTMLURL string `json:"html_url"` + PublishedAt string `json:"published_at"` + CreatedAt string `json:"created_at"` + Draft bool `json:"draft"` + Prerelease bool `json:"prerelease"` +} + +type changelogReleaseMetadata struct { + Provider string `json:"provider"` + Repository string `json:"repository,omitempty"` + TagName string `json:"tag_name,omitempty"` + Prerelease bool `json:"prerelease,omitempty"` + RawPublishedAt string `json:"raw_published_at,omitempty"` + RawCreatedAt string `json:"raw_created_at,omitempty"` +} + +type changelogRSSFeed struct { + XMLName xml.Name `xml:"rss"` + Channel changelogRSSChannel `xml:"channel"` +} + +type changelogRSSChannel struct { + Title string `xml:"title"` + Items []changelogRSSFeedItem `xml:"item"` +} + +type changelogRSSFeedItem struct { + Title string `xml:"title"` + Link string `xml:"link"` + Description string `xml:"description"` + ContentEncoded string `xml:"encoded"` + GUID changelogRSSFeedGUID `xml:"guid"` + PubDate string `xml:"pubDate"` +} + +type changelogRSSFeedGUID struct { + Value string `xml:",chardata"` +} + +type changelogEntryMetadata struct { + Provider string `json:"provider"` + FeedURL string `json:"feed_url,omitempty"` + FeedTitle string `json:"feed_title,omitempty"` + GUID string `json:"guid,omitempty"` + RawPubDate string `json:"raw_pub_date,omitempty"` +} + +func NewChangelogFetcher(options ChangelogFetcherOptions) *ChangelogFetcher { + githubHTTPClient := options.GitHubHTTPClient + if githubHTTPClient == nil { + githubHTTPClient = http.DefaultClient + } + + githubUserAgent := strings.TrimSpace(options.GitHubUserAgent) + if githubUserAgent == "" { + githubUserAgent = defaultGitHubUserAgent + } + + githubToken := strings.TrimSpace(options.GitHubToken) + if githubToken == "" { + githubToken = strings.TrimSpace(os.Getenv(githubTokenEnvVar)) + } + + rssHTTPClient := options.RSSHTTPClient + if rssHTTPClient == nil { + rssHTTPClient = http.DefaultClient + } + + rssUserAgent := strings.TrimSpace(options.RSSUserAgent) + if rssUserAgent == "" { + rssUserAgent = defaultChangelogRSSUserAgent + } + + now := options.Now + if now == nil { + now = time.Now + } + + return &ChangelogFetcher{ + githubBaseURL: mustParseChangelogBaseURL(normalizeChangelogBaseURL(options.GitHubBaseURL, defaultGitHubAPIBaseURL)), + githubHTTPClient: githubHTTPClient, + githubUserAgent: githubUserAgent, + githubToken: githubToken, + rssHTTPClient: rssHTTPClient, + rssUserAgent: rssUserAgent, + now: now, + } +} + +func (fetcher *ChangelogFetcher) Fetch(ctx context.Context, request Request) (Result, error) { + ctx = normalizeContext(ctx) + + if request.DB == nil { + return Result{}, fmt.Errorf("database handle is required") + } + if request.FetchRun.ID <= 0 { + return Result{}, fmt.Errorf("fetch run is required") + } + if strings.TrimSpace(request.DatasetID) == "" { + return Result{}, fmt.Errorf("dataset ID is required") + } + if strings.TrimSpace(request.Binding.ID) == "" { + return Result{}, fmt.Errorf("binding ID is required") + } + if strings.TrimSpace(request.Binding.EntityID) == "" { + return Result{}, fmt.Errorf("binding %q entity_id is required", request.Binding.ID) + } + if strings.TrimSpace(request.Source.ID) == "" { + return Result{}, fmt.Errorf("source ID is required") + } + + scope, err := parseChangelogBindingScope(request.Binding) + if err != nil { + return Result{}, err + } + + discoveredAt := request.FetchRun.StartedAt.UTC() + if discoveredAt.IsZero() { + discoveredAt = fetcher.now().UTC() + } + + var contentItems []storage.ContentItemInput + + switch scope.Mode { + case changelogModeGitHubReleases: + contentItems, err = fetcher.fetchGitHubReleaseContentItems(ctx, scope.Repositories, discoveredAt) + if err != nil { + return Result{}, fmt.Errorf("binding %q GitHub releases: %w", request.Binding.ID, err) + } + case changelogModeRSS: + contentItems, err = fetcher.fetchRSSChangelogContentItems(ctx, scope.FeedURL, discoveredAt) + if err != nil { + return Result{}, fmt.Errorf("binding %q RSS changelog: %w", request.Binding.ID, err) + } + default: + return Result{}, fmt.Errorf("binding %q has unsupported changelog mode %q", request.Binding.ID, scope.Mode) + } + + contentItemsInserted, err := storage.InsertContentItems(ctx, request.DB, request.FetchRun.ID, contentItems) + if err != nil { + return Result{}, fmt.Errorf("persist changelog content items: %w", err) + } + + return Result{ + RecordsWritten: len(contentItems), + MetricsWritten: 0, + ContentItemsWritten: contentItemsInserted, + }, nil +} + +func parseChangelogBindingScope(binding config.Binding) (changelogBindingScope, error) { + rawMode, exists := binding.Scope["mode"] + if !exists { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.mode is required", binding.ID) + } + + modeValue, ok := rawMode.(string) + if !ok { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.mode must be a string", binding.ID) + } + + mode := changelogBindingMode(strings.TrimSpace(modeValue)) + switch mode { + case changelogModeGitHubReleases: + rawRepos, exists := binding.Scope["repos"] + if !exists { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.repos is required when mode is github_releases", binding.ID) + } + + repositories, err := parseChangelogRepositoryRefs(binding.ID, rawRepos) + if err != nil { + return changelogBindingScope{}, err + } + + return changelogBindingScope{ + Mode: mode, + Repositories: repositories, + }, nil + + case changelogModeRSS: + rawFeedURL, exists := binding.Scope["feed_url"] + if !exists { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.feed_url is required when mode is rss", binding.ID) + } + + feedURLValue, ok := rawFeedURL.(string) + if !ok { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.feed_url must be a string", binding.ID) + } + + feedURLValue = strings.TrimSpace(feedURLValue) + if feedURLValue == "" { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.feed_url must not be empty", binding.ID) + } + + parsedFeedURL, err := url.Parse(feedURLValue) + if err != nil { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.feed_url parse error: %w", binding.ID, err) + } + if !parsedFeedURL.IsAbs() { + return changelogBindingScope{}, fmt.Errorf("binding %q scope.feed_url must be an absolute URL", binding.ID) + } + + return changelogBindingScope{ + Mode: mode, + FeedURL: parsedFeedURL, + }, nil + + default: + return changelogBindingScope{}, fmt.Errorf("binding %q has unsupported changelog mode %q", binding.ID, mode) + } +} + +func parseChangelogRepositoryRefs(bindingID string, value any) ([]githubRepositoryRef, error) { + var rawValues []string + + switch typed := value.(type) { + case []string: + rawValues = append(rawValues, typed...) + case []any: + rawValues = make([]string, 0, len(typed)) + for index, item := range typed { + text, ok := item.(string) + if !ok { + return nil, fmt.Errorf("binding %q scope.repos[%d] must be a string", bindingID, index) + } + rawValues = append(rawValues, text) + } + default: + return nil, fmt.Errorf("binding %q scope.repos must be an array of repository strings", bindingID) + } + + if len(rawValues) == 0 { + return nil, fmt.Errorf("binding %q scope.repos must not be empty", bindingID) + } + + repositories := make([]githubRepositoryRef, 0, len(rawValues)) + seen := make(map[string]struct{}, len(rawValues)) + + for index, rawRepository := range rawValues { + repositoryRef, err := parseGitHubRepositoryRef(strings.TrimSpace(rawRepository)) + if err != nil { + return nil, fmt.Errorf("binding %q scope.repos[%d]: %w", bindingID, index, err) + } + + fullName := repositoryRef.FullName() + if _, exists := seen[fullName]; exists { + continue + } + + seen[fullName] = struct{}{} + repositories = append(repositories, repositoryRef) + } + + if len(repositories) == 0 { + return nil, fmt.Errorf("binding %q scope.repos must yield at least one repository", bindingID) + } + + return repositories, nil +} + +func (fetcher *ChangelogFetcher) fetchGitHubReleaseContentItems( + ctx context.Context, + repositories []githubRepositoryRef, + discoveredAt time.Time, +) ([]storage.ContentItemInput, error) { + normalized := make([]storage.ContentItemInput, 0) + seen := make(map[string]struct{}) + + for _, repositoryRef := range repositories { + releases, err := fetcher.listGitHubReleases(ctx, repositoryRef) + if err != nil { + return nil, fmt.Errorf("repo %q: %w", repositoryRef.FullName(), err) + } + + for _, release := range releases { + normalizedItem, skip, err := normalizeGitHubReleaseContentItem(repositoryRef, release, discoveredAt) + if err != nil { + return nil, err + } + if skip { + continue + } + if _, exists := seen[normalizedItem.DedupeKey]; exists { + continue + } + + seen[normalizedItem.DedupeKey] = struct{}{} + normalized = append(normalized, normalizedItem) + } + } + + return normalized, nil +} + +func (fetcher *ChangelogFetcher) listGitHubReleases( + ctx context.Context, + repositoryRef githubRepositoryRef, +) ([]changelogGitHubRelease, error) { + releases := make([]changelogGitHubRelease, 0) + + for page := 1; page <= changelogGitHubMaxPages; page++ { + query := url.Values{} + query.Set("page", strconv.Itoa(page)) + query.Set("per_page", strconv.Itoa(changelogGitHubReleasesPageSize)) + + response, err := fetcher.doGitHubRequest(ctx, repositoryPath(repositoryRef)+"/releases", query) + if err != nil { + return nil, err + } + + pageReleases, err := fetcher.decodeGitHubReleasePage(response) + if err != nil { + return nil, err + } + + if len(pageReleases) == 0 { + return releases, nil + } + + releases = append(releases, pageReleases...) + } + + return nil, fmt.Errorf( + "repository %q exceeded %d GitHub release pages", + repositoryRef.FullName(), + changelogGitHubMaxPages, + ) +} + +func (fetcher *ChangelogFetcher) decodeGitHubReleasePage(response *http.Response) ([]changelogGitHubRelease, error) { + defer response.Body.Close() + + if response.StatusCode != http.StatusOK { + return nil, fetcher.githubResponseError("list releases", response) + } + + var releases []changelogGitHubRelease + if err := json.NewDecoder(response.Body).Decode(&releases); err != nil { + return nil, fmt.Errorf("decode GitHub releases response: %w", err) + } + + return releases, nil +} + +func normalizeGitHubReleaseContentItem( + repositoryRef githubRepositoryRef, + release changelogGitHubRelease, + discoveredAt time.Time, +) (storage.ContentItemInput, bool, error) { + if release.Draft { + return storage.ContentItemInput{}, true, nil + } + + publishedAt, rawPublishedAt := parseChangelogRFC3339Timestamp(release.PublishedAt) + if publishedAt == nil { + publishedAt, _ = parseChangelogRFC3339Timestamp(release.CreatedAt) + } + + metadataJSON, err := buildGitHubReleaseMetadata(repositoryRef, release, rawPublishedAt) + if err != nil { + return storage.ContentItemInput{}, false, err + } + + title := strings.TrimSpace(release.Name) + if title == "" { + title = strings.TrimSpace(release.TagName) + } + + input := storage.ContentItemInput{ + ItemType: ContentItemTypeRelease, + Title: title, + Summary: strings.TrimSpace(release.Body), + URL: strings.TrimSpace(release.HTMLURL), + ExternalID: githubReleaseExternalID(release.ID), + PublishedAt: publishedAt, + DiscoveredAt: discoveredAt, + MetadataJSON: metadataJSON, + } + + dedupeKey, err := storage.BuildContentItemDedupeKey(input) + if err != nil { + return storage.ContentItemInput{}, true, nil + } + + input.DedupeKey = dedupeKey + + return input, false, nil +} + +func buildGitHubReleaseMetadata( + repositoryRef githubRepositoryRef, + release changelogGitHubRelease, + rawPublishedAt string, +) (string, error) { + metadata := changelogReleaseMetadata{ + Provider: changelogProviderGitHubReleases, + Repository: repositoryRef.FullName(), + TagName: strings.TrimSpace(release.TagName), + Prerelease: release.Prerelease, + RawPublishedAt: strings.TrimSpace(rawPublishedAt), + RawCreatedAt: strings.TrimSpace(release.CreatedAt), + } + + encoded, err := json.Marshal(metadata) + if err != nil { + return "", fmt.Errorf("marshal GitHub release metadata: %w", err) + } + + return string(encoded), nil +} + +func githubReleaseExternalID(id int64) string { + if id <= 0 { + return "" + } + + return strconv.FormatInt(id, 10) +} + +func parseChangelogRFC3339Timestamp(value string) (*time.Time, string) { + raw := strings.TrimSpace(value) + if raw == "" { + return nil, "" + } + + parsed, err := time.Parse(time.RFC3339, raw) + if err != nil { + return nil, raw + } + + timestamp := parsed.UTC() + return ×tamp, raw +} + +func (fetcher *ChangelogFetcher) doGitHubRequest( + ctx context.Context, + path string, + query url.Values, +) (*http.Response, error) { + path = strings.TrimPrefix(strings.TrimSpace(path), "/") + + requestURL := fetcher.githubBaseURL.ResolveReference(&url.URL{ + Path: path, + RawQuery: query.Encode(), + }).String() + + request, err := http.NewRequestWithContext(ctx, http.MethodGet, requestURL, nil) + if err != nil { + return nil, fmt.Errorf("build changelog GitHub request: %w", err) + } + + request.Header.Set("Accept", "application/vnd.github+json") + request.Header.Set("User-Agent", fetcher.githubUserAgent) + if fetcher.githubToken != "" { + request.Header.Set("Authorization", "Bearer "+fetcher.githubToken) + } + + response, err := fetcher.githubHTTPClient.Do(request) + if err != nil { + return nil, fmt.Errorf("perform changelog GitHub request %s: %w", requestURL, err) + } + + return response, nil +} + +func (fetcher *ChangelogFetcher) githubResponseError(action string, response *http.Response) error { + body, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + + message := strings.TrimSpace(string(body)) + if message == "" { + message = http.StatusText(response.StatusCode) + } + if message == "" { + message = "unknown GitHub changelog error" + } + + requestPath := "" + if response.Request != nil && response.Request.URL != nil { + requestPath = response.Request.URL.Path + } + + return fmt.Errorf("GitHub %s %s returned %d: %s", action, requestPath, response.StatusCode, message) +} + +func (fetcher *ChangelogFetcher) fetchRSSChangelogContentItems( + ctx context.Context, + feedURL *url.URL, + discoveredAt time.Time, +) ([]storage.ContentItemInput, error) { + response, err := fetcher.doRSSRequest(ctx, feedURL.String()) + if err != nil { + return nil, err + } + defer response.Body.Close() + + if response.StatusCode != http.StatusOK { + return nil, fetcher.rssResponseError(response) + } + + var feed changelogRSSFeed + if err := xml.NewDecoder(response.Body).Decode(&feed); err != nil { + return nil, fmt.Errorf("decode changelog RSS response: %w", err) + } + + return normalizeRSSChangelogContentItems(feedURL, strings.TrimSpace(feed.Channel.Title), feed.Channel.Items, discoveredAt) +} + +func normalizeRSSChangelogContentItems( + feedURL *url.URL, + feedTitle string, + items []changelogRSSFeedItem, + discoveredAt time.Time, +) ([]storage.ContentItemInput, error) { + normalized := make([]storage.ContentItemInput, 0, len(items)) + seen := make(map[string]struct{}, len(items)) + + for _, item := range items { + normalizedItem, skip, err := normalizeRSSChangelogContentItem(feedURL, feedTitle, item, discoveredAt) + if err != nil { + return nil, err + } + if skip { + continue + } + if _, exists := seen[normalizedItem.DedupeKey]; exists { + continue + } + + seen[normalizedItem.DedupeKey] = struct{}{} + normalized = append(normalized, normalizedItem) + } + + return normalized, nil +} + +func normalizeRSSChangelogContentItem( + feedURL *url.URL, + feedTitle string, + item changelogRSSFeedItem, + discoveredAt time.Time, +) (storage.ContentItemInput, bool, error) { + publishedAt, rawPubDate := parseNewsRSSPublishedAt(item.PubDate) + + summary := strings.TrimSpace(item.ContentEncoded) + if summary == "" { + summary = strings.TrimSpace(item.Description) + } + + metadataJSON, err := buildRSSChangelogMetadata(feedURL.String(), feedTitle, item, rawPubDate) + if err != nil { + return storage.ContentItemInput{}, false, err + } + + input := storage.ContentItemInput{ + ItemType: ContentItemTypeChangelogEntry, + Title: strings.TrimSpace(item.Title), + Summary: summary, + URL: resolveNewsRSSLink(feedURL, item.Link), + ExternalID: strings.TrimSpace(item.GUID.Value), + PublishedAt: publishedAt, + DiscoveredAt: discoveredAt, + MetadataJSON: metadataJSON, + } + + dedupeKey, err := storage.BuildContentItemDedupeKey(input) + if err != nil { + return storage.ContentItemInput{}, true, nil + } + + input.DedupeKey = dedupeKey + + return input, false, nil +} + +func buildRSSChangelogMetadata( + feedURL string, + feedTitle string, + item changelogRSSFeedItem, + rawPubDate string, +) (string, error) { + metadata := changelogEntryMetadata{ + Provider: changelogProviderRSS, + FeedURL: strings.TrimSpace(feedURL), + FeedTitle: strings.TrimSpace(feedTitle), + GUID: strings.TrimSpace(item.GUID.Value), + RawPubDate: strings.TrimSpace(rawPubDate), + } + + encoded, err := json.Marshal(metadata) + if err != nil { + return "", fmt.Errorf("marshal changelog RSS metadata: %w", err) + } + + return string(encoded), nil +} + +func (fetcher *ChangelogFetcher) doRSSRequest(ctx context.Context, feedURL string) (*http.Response, error) { + request, err := http.NewRequestWithContext(ctx, http.MethodGet, feedURL, nil) + if err != nil { + return nil, fmt.Errorf("build changelog RSS request: %w", err) + } + + request.Header.Set("Accept", "application/rss+xml, application/xml;q=0.9, text/xml;q=0.8") + request.Header.Set("User-Agent", fetcher.rssUserAgent) + + response, err := fetcher.rssHTTPClient.Do(request) + if err != nil { + return nil, fmt.Errorf("perform changelog RSS request %s: %w", feedURL, err) + } + + return response, nil +} + +func (fetcher *ChangelogFetcher) rssResponseError(response *http.Response) error { + body, _ := io.ReadAll(io.LimitReader(response.Body, 4096)) + + message := strings.TrimSpace(string(body)) + if message == "" { + message = http.StatusText(response.StatusCode) + } + if message == "" { + message = "unknown changelog RSS error" + } + + requestPath := "" + if response.Request != nil && response.Request.URL != nil { + requestPath = response.Request.URL.Path + } + + return fmt.Errorf("changelog RSS GET %s returned %d: %s", requestPath, response.StatusCode, message) +} + +func normalizeChangelogBaseURL(value string, fallback string) string { + trimmed := strings.TrimSpace(value) + if trimmed == "" { + trimmed = fallback + } + if !strings.HasSuffix(trimmed, "/") { + trimmed += "/" + } + + return trimmed +} + +func mustParseChangelogBaseURL(value string) *url.URL { + parsed, err := url.Parse(value) + if err != nil { + panic(fmt.Errorf("parse changelog base URL %q: %w", value, err)) + } + + return parsed +} diff --git a/internal/fetch/changelog_test.go b/internal/fetch/changelog_test.go new file mode 100644 index 0000000..b381d90 --- /dev/null +++ b/internal/fetch/changelog_test.go @@ -0,0 +1,839 @@ +package fetch + +import ( + "context" + "database/sql" + "io" + "net/http" + "net/http/httptest" + "strings" + "sync" + "testing" + "time" + + "github.com/CrisSTEM/signalscope/internal/config" + "github.com/CrisSTEM/signalscope/internal/storage" +) + +type changelogContentItemRow struct { + BindingID string + ItemType string + Title string + Summary string + URL string + ExternalID string + PublishedAt sql.NullString + DiscoveredAt string + MetadataJSON string +} + +type changelogHTTPResponse struct { + Status int + Body string + ContentType string +} + +type changelogHTTPTestServerOptions struct { + Responses map[string][]changelogHTTPResponse +} + +func TestDefaultRegistryRegistersChangelogFetcher(t *testing.T) { + t.Parallel() + + registry := DefaultRegistry() + if _, ok := registry.Lookup(SourceKindChangelog); !ok { + t.Fatalf("DefaultRegistry().Lookup(%q) ok = false, want true", SourceKindChangelog) + } +} + +func TestChangelogFetcherGitHubReleasesPersistsReleaseItems(t *testing.T) { + t.Parallel() + + ctx := context.Background() + discoveredAt := time.Date(2026, 4, 22, 10, 0, 0, 0, time.UTC) + + server := newChangelogHTTPTestServer(t, changelogHTTPTestServerOptions{ + Responses: map[string][]changelogHTTPResponse{ + "/repos/example/repo/releases?page=1&per_page=100": { + { + Body: `[ + { + "id": 101, + "tag_name": "v2.0.0", + "name": "Release 2.0.0", + "body": "Major release body", + "html_url": "https://github.com/example/repo/releases/tag/v2.0.0", + "published_at": "2026-04-21T09:30:00Z", + "created_at": "2026-04-20T09:30:00Z" + }, + { + "id": 102, + "tag_name": "v1.9.0", + "name": "", + "body": "Fallback title should use tag name", + "html_url": "https://github.com/example/repo/releases/tag/v1.9.0", + "published_at": "2026-04-20T09:30:00Z", + "created_at": "2026-04-19T09:30:00Z", + "prerelease": true + } +]`, + }, + }, + "/repos/example/repo/releases?page=2&per_page=100": { + { + Body: `[ + { + "id": 103, + "tag_name": "v1.8.1-draft", + "name": "Draft release should be skipped", + "body": "Draft body", + "html_url": "https://github.com/example/repo/releases/tag/v1.8.1-draft", + "published_at": "2026-04-19T09:30:00Z", + "created_at": "2026-04-18T09:30:00Z", + "draft": true + }, + { + "id": 104, + "tag_name": "v1.8.0", + "name": "Release 1.8.0", + "body": "Older stable release", + "html_url": "https://github.com/example/repo/releases/tag/v1.8.0", + "published_at": "2026-04-18T09:30:00Z", + "created_at": "2026-04-17T09:30:00Z" + } +]`, + }, + }, + "/repos/example/repo/releases?page=3&per_page=100": { + { + Body: `[]`, + }, + }, + }, + }) + defer server.Close() + + pack := runtimeFixturePackWithChangelogBinding( + "changelog-github-success", + "changelog-binding-github-1", + config.JSONMap{ + "mode": "github_releases", + "repos": []string{"example/repo"}, + }, + ) + + db := openRuntimeTestDB(t, pack) + defer db.Close() + + registry := NewRegistry() + registry.MustRegister(SourceKindChangelog, newChangelogTestFetcher(server)) + + runtime := Runtime{ + DB: db, + Registry: registry, + Now: sequenceClock( + discoveredAt, + discoveredAt.Add(time.Second), + ), + } + + summary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-github-1", + }) + if err != nil { + t.Fatalf("Execute() error = %v", err) + } + + if summary.Attempted != 1 { + t.Fatalf("summary.Attempted = %d, want %d", summary.Attempted, 1) + } + if summary.Succeeded != 1 { + t.Fatalf("summary.Succeeded = %d, want %d", summary.Succeeded, 1) + } + if summary.Failed != 0 { + t.Fatalf("summary.Failed = %d, want %d", summary.Failed, 0) + } + if summary.RecordsWritten != 3 { + t.Fatalf("summary.RecordsWritten = %d, want %d", summary.RecordsWritten, 3) + } + if summary.MetricsWritten != 0 { + t.Fatalf("summary.MetricsWritten = %d, want %d", summary.MetricsWritten, 0) + } + if summary.ContentItemsWritten != 3 { + t.Fatalf("summary.ContentItemsWritten = %d, want %d", summary.ContentItemsWritten, 3) + } + if len(summary.Results) != 1 { + t.Fatalf("len(summary.Results) = %d, want %d", len(summary.Results), 1) + } + if summary.Results[0].Status != storage.FetchRunStatusSucceeded { + t.Fatalf("summary.Results[0].Status = %q, want %q", summary.Results[0].Status, storage.FetchRunStatusSucceeded) + } + + rows, err := readChangelogContentItemsForBinding(ctx, db, "changelog-binding-github-1") + if err != nil { + t.Fatalf("readChangelogContentItemsForBinding() error = %v", err) + } + if len(rows) != 3 { + t.Fatalf("len(rows) = %d, want %d", len(rows), 3) + } + + if rows[0].ItemType != ContentItemTypeRelease { + t.Fatalf("rows[0].ItemType = %q, want %q", rows[0].ItemType, ContentItemTypeRelease) + } + if rows[0].Title != "Release 2.0.0" { + t.Fatalf("rows[0].Title = %q, want %q", rows[0].Title, "Release 2.0.0") + } + if rows[0].Summary != "Major release body" { + t.Fatalf("rows[0].Summary = %q, want %q", rows[0].Summary, "Major release body") + } + if rows[0].URL != "https://github.com/example/repo/releases/tag/v2.0.0" { + t.Fatalf("rows[0].URL = %q, want release URL", rows[0].URL) + } + if rows[0].ExternalID != "101" { + t.Fatalf("rows[0].ExternalID = %q, want %q", rows[0].ExternalID, "101") + } + if !rows[0].PublishedAt.Valid || rows[0].PublishedAt.String != "2026-04-21T09:30:00Z" { + t.Fatalf("rows[0].PublishedAt = %+v, want %q", rows[0].PublishedAt, "2026-04-21T09:30:00Z") + } + if rows[0].DiscoveredAt != discoveredAt.Format(time.RFC3339) { + t.Fatalf("rows[0].DiscoveredAt = %q, want %q", rows[0].DiscoveredAt, discoveredAt.Format(time.RFC3339)) + } + if !strings.Contains(rows[0].MetadataJSON, `"provider":"github_releases"`) { + t.Fatalf("rows[0].MetadataJSON = %q, want provider marker", rows[0].MetadataJSON) + } + if !strings.Contains(rows[0].MetadataJSON, `"repository":"example/repo"`) { + t.Fatalf("rows[0].MetadataJSON = %q, want repository marker", rows[0].MetadataJSON) + } + if !strings.Contains(rows[0].MetadataJSON, `"tag_name":"v2.0.0"`) { + t.Fatalf("rows[0].MetadataJSON = %q, want tag marker", rows[0].MetadataJSON) + } + + if rows[1].Title != "v1.9.0" { + t.Fatalf("rows[1].Title = %q, want tag-name fallback", rows[1].Title) + } + if rows[1].ExternalID != "102" { + t.Fatalf("rows[1].ExternalID = %q, want %q", rows[1].ExternalID, "102") + } + + if rows[2].Title != "Release 1.8.0" { + t.Fatalf("rows[2].Title = %q, want %q", rows[2].Title, "Release 1.8.0") + } + + for _, row := range rows { + if row.Title == "Draft release should be skipped" { + t.Fatal("unexpected draft release persisted") + } + } +} + +func TestChangelogFetcherGitHubReleasesRerunDedupesExistingItems(t *testing.T) { + t.Parallel() + + ctx := context.Background() + firstDiscoveredAt := time.Date(2026, 4, 22, 11, 0, 0, 0, time.UTC) + secondDiscoveredAt := time.Date(2026, 4, 22, 12, 0, 0, 0, time.UTC) + + server := newChangelogHTTPTestServer(t, changelogHTTPTestServerOptions{ + Responses: map[string][]changelogHTTPResponse{ + "/repos/example/repo/releases?page=1&per_page=100": { + { + Body: `[ + { + "id": 201, + "tag_name": "v1.0.0", + "name": "Release 1.0.0", + "body": "Original release body", + "html_url": "https://github.com/example/repo/releases/tag/v1.0.0", + "published_at": "2026-04-21T10:00:00Z", + "created_at": "2026-04-20T10:00:00Z" + }, + { + "id": 202, + "tag_name": "v0.9.0", + "name": "Release 0.9.0", + "body": "Original second release body", + "html_url": "https://github.com/example/repo/releases/tag/v0.9.0", + "published_at": "2026-04-20T10:00:00Z", + "created_at": "2026-04-19T10:00:00Z" + } +]`, + }, + { + Body: `[ + { + "id": 201, + "tag_name": "v1.0.0", + "name": "Release 1.0.0 UPDATED", + "body": "Updated body should not overwrite the first row", + "html_url": "https://github.com/example/repo/releases/tag/v1.0.0", + "published_at": "2026-04-21T10:00:00Z", + "created_at": "2026-04-20T10:00:00Z" + }, + { + "id": 202, + "tag_name": "v0.9.0", + "name": "Release 0.9.0 UPDATED", + "body": "Updated second body should not overwrite the first row", + "html_url": "https://github.com/example/repo/releases/tag/v0.9.0", + "published_at": "2026-04-20T10:00:00Z", + "created_at": "2026-04-19T10:00:00Z" + } +]`, + }, + }, + "/repos/example/repo/releases?page=2&per_page=100": { + {Body: `[]`}, + {Body: `[]`}, + }, + }, + }) + defer server.Close() + + pack := runtimeFixturePackWithChangelogBinding( + "changelog-github-rerun", + "changelog-binding-github-1", + config.JSONMap{ + "mode": "github_releases", + "repos": []string{"example/repo"}, + }, + ) + + db := openRuntimeTestDB(t, pack) + defer db.Close() + + registry := NewRegistry() + registry.MustRegister(SourceKindChangelog, newChangelogTestFetcher(server)) + + runtime := Runtime{ + DB: db, + Registry: registry, + Now: sequenceClock( + firstDiscoveredAt, + firstDiscoveredAt.Add(time.Second), + secondDiscoveredAt, + secondDiscoveredAt.Add(time.Second), + ), + } + + firstSummary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-github-1", + }) + if err != nil { + t.Fatalf("Execute(first run) error = %v", err) + } + if firstSummary.ContentItemsWritten != 2 { + t.Fatalf("firstSummary.ContentItemsWritten = %d, want %d", firstSummary.ContentItemsWritten, 2) + } + if firstSummary.RecordsWritten != 2 { + t.Fatalf("firstSummary.RecordsWritten = %d, want %d", firstSummary.RecordsWritten, 2) + } + + secondSummary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-github-1", + }) + if err != nil { + t.Fatalf("Execute(second run) error = %v", err) + } + if secondSummary.ContentItemsWritten != 0 { + t.Fatalf("secondSummary.ContentItemsWritten = %d, want %d", secondSummary.ContentItemsWritten, 0) + } + if secondSummary.RecordsWritten != 2 { + t.Fatalf("secondSummary.RecordsWritten = %d, want %d", secondSummary.RecordsWritten, 2) + } + + contentCount, err := queryCount(ctx, db, `SELECT COUNT(*) FROM content_items WHERE binding_id = ?`, "changelog-binding-github-1") + if err != nil { + t.Fatalf("queryCount(content_items) error = %v", err) + } + if contentCount != 2 { + t.Fatalf("contentCount = %d, want %d", contentCount, 2) + } + + rows, err := readChangelogContentItemsForBinding(ctx, db, "changelog-binding-github-1") + if err != nil { + t.Fatalf("readChangelogContentItemsForBinding() error = %v", err) + } + if len(rows) != 2 { + t.Fatalf("len(rows) = %d, want %d", len(rows), 2) + } + if rows[0].Title != "Release 1.0.0" { + t.Fatalf("rows[0].Title = %q, want first-write title", rows[0].Title) + } + if rows[0].Summary != "Original release body" { + t.Fatalf("rows[0].Summary = %q, want first-write summary", rows[0].Summary) + } + if rows[0].DiscoveredAt != firstDiscoveredAt.Format(time.RFC3339) { + t.Fatalf("rows[0].DiscoveredAt = %q, want %q", rows[0].DiscoveredAt, firstDiscoveredAt.Format(time.RFC3339)) + } +} + +func TestChangelogFetcherRSSPersistsEntriesAndRerunDedupes(t *testing.T) { + t.Parallel() + + ctx := context.Background() + firstDiscoveredAt := time.Date(2026, 4, 22, 13, 0, 0, 0, time.UTC) + secondDiscoveredAt := time.Date(2026, 4, 22, 14, 0, 0, 0, time.UTC) + + server := newChangelogHTTPTestServer(t, changelogHTTPTestServerOptions{ + Responses: map[string][]changelogHTTPResponse{ + "/feeds/changelog.xml": { + { + ContentType: "application/rss+xml; charset=utf-8", + Body: ` + + + Example Changelog + + Version 2.0 shipped + https://example.com/changelog/v2 + Original changelog body + guid-1 + Tue, 22 Apr 2026 12:30:00 GMT + + + Version 2.0 shipped duplicate + https://example.com/changelog/v2 + Duplicate entry should be ignored + guid-1 + Tue, 22 Apr 2026 12:30:00 GMT + + + Migration note published + /posts/migration-note + Second changelog entry + Tue, 22 Apr 2026 12:45:00 GMT + + +`, + }, + { + ContentType: "application/rss+xml; charset=utf-8", + Body: ` + + + Example Changelog + + Version 2.0 shipped UPDATED + https://example.com/changelog/v2 + Updated body should not overwrite first row + guid-1 + Tue, 22 Apr 2026 12:30:00 GMT + + + Migration note published UPDATED + /posts/migration-note + Updated second entry should not overwrite first row + Tue, 22 Apr 2026 12:45:00 GMT + + +`, + }, + }, + }, + }) + defer server.Close() + + pack := runtimeFixturePackWithChangelogBinding( + "changelog-rss-rerun", + "changelog-binding-rss-1", + config.JSONMap{ + "mode": "rss", + "feed_url": server.URL + "/feeds/changelog.xml", + }, + ) + + db := openRuntimeTestDB(t, pack) + defer db.Close() + + registry := NewRegistry() + registry.MustRegister(SourceKindChangelog, newChangelogTestFetcher(server)) + + runtime := Runtime{ + DB: db, + Registry: registry, + Now: sequenceClock( + firstDiscoveredAt, + firstDiscoveredAt.Add(time.Second), + secondDiscoveredAt, + secondDiscoveredAt.Add(time.Second), + ), + } + + firstSummary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-rss-1", + }) + if err != nil { + t.Fatalf("Execute(first run) error = %v", err) + } + if firstSummary.ContentItemsWritten != 2 { + t.Fatalf("firstSummary.ContentItemsWritten = %d, want %d", firstSummary.ContentItemsWritten, 2) + } + if firstSummary.RecordsWritten != 2 { + t.Fatalf("firstSummary.RecordsWritten = %d, want %d", firstSummary.RecordsWritten, 2) + } + + secondSummary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-rss-1", + }) + if err != nil { + t.Fatalf("Execute(second run) error = %v", err) + } + if secondSummary.ContentItemsWritten != 0 { + t.Fatalf("secondSummary.ContentItemsWritten = %d, want %d", secondSummary.ContentItemsWritten, 0) + } + if secondSummary.RecordsWritten != 2 { + t.Fatalf("secondSummary.RecordsWritten = %d, want %d", secondSummary.RecordsWritten, 2) + } + + rows, err := readChangelogContentItemsForBinding(ctx, db, "changelog-binding-rss-1") + if err != nil { + t.Fatalf("readChangelogContentItemsForBinding() error = %v", err) + } + if len(rows) != 2 { + t.Fatalf("len(rows) = %d, want %d", len(rows), 2) + } + + if rows[0].ItemType != ContentItemTypeChangelogEntry { + t.Fatalf("rows[0].ItemType = %q, want %q", rows[0].ItemType, ContentItemTypeChangelogEntry) + } + if rows[0].Title != "Version 2.0 shipped" { + t.Fatalf("rows[0].Title = %q, want %q", rows[0].Title, "Version 2.0 shipped") + } + if rows[0].Summary != "Original changelog body" { + t.Fatalf("rows[0].Summary = %q, want %q", rows[0].Summary, "Original changelog body") + } + if rows[0].URL != "https://example.com/changelog/v2" { + t.Fatalf("rows[0].URL = %q, want %q", rows[0].URL, "https://example.com/changelog/v2") + } + if rows[0].ExternalID != "guid-1" { + t.Fatalf("rows[0].ExternalID = %q, want %q", rows[0].ExternalID, "guid-1") + } + if !rows[0].PublishedAt.Valid || rows[0].PublishedAt.String != "2026-04-22T12:30:00Z" { + t.Fatalf("rows[0].PublishedAt = %+v, want %q", rows[0].PublishedAt, "2026-04-22T12:30:00Z") + } + if rows[0].DiscoveredAt != firstDiscoveredAt.Format(time.RFC3339) { + t.Fatalf("rows[0].DiscoveredAt = %q, want %q", rows[0].DiscoveredAt, firstDiscoveredAt.Format(time.RFC3339)) + } + if !strings.Contains(rows[0].MetadataJSON, `"provider":"rss"`) { + t.Fatalf("rows[0].MetadataJSON = %q, want provider marker", rows[0].MetadataJSON) + } + if !strings.Contains(rows[0].MetadataJSON, `"feed_url":"`+server.URL+`/feeds/changelog.xml"`) { + t.Fatalf("rows[0].MetadataJSON = %q, want feed URL marker", rows[0].MetadataJSON) + } + if !strings.Contains(rows[0].MetadataJSON, `"guid":"guid-1"`) { + t.Fatalf("rows[0].MetadataJSON = %q, want guid marker", rows[0].MetadataJSON) + } + + if rows[1].URL != server.URL+"/posts/migration-note" { + t.Fatalf("rows[1].URL = %q, want resolved absolute URL", rows[1].URL) + } + if rows[1].Title != "Migration note published" { + t.Fatalf("rows[1].Title = %q, want %q", rows[1].Title, "Migration note published") + } + if rows[1].Summary != "Second changelog entry" { + t.Fatalf("rows[1].Summary = %q, want %q", rows[1].Summary, "Second changelog entry") + } +} + +func TestChangelogFetcherFailsOnMalformedGitHubReleasesResponse(t *testing.T) { + t.Parallel() + + ctx := context.Background() + discoveredAt := time.Date(2026, 4, 22, 15, 0, 0, 0, time.UTC) + + server := newChangelogHTTPTestServer(t, changelogHTTPTestServerOptions{ + Responses: map[string][]changelogHTTPResponse{ + "/repos/example/repo/releases?page=1&per_page=100": { + { + Body: `[`, + }, + }, + }, + }) + defer server.Close() + + pack := runtimeFixturePackWithChangelogBinding( + "changelog-github-malformed", + "changelog-binding-github-1", + config.JSONMap{ + "mode": "github_releases", + "repos": []string{"example/repo"}, + }, + ) + + db := openRuntimeTestDB(t, pack) + defer db.Close() + + registry := NewRegistry() + registry.MustRegister(SourceKindChangelog, newChangelogTestFetcher(server)) + + runtime := Runtime{ + DB: db, + Registry: registry, + Now: sequenceClock( + discoveredAt, + discoveredAt.Add(time.Second), + ), + } + + summary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-github-1", + }) + if err != nil { + t.Fatalf("Execute() error = %v", err) + } + + if summary.Attempted != 1 { + t.Fatalf("summary.Attempted = %d, want %d", summary.Attempted, 1) + } + if summary.Succeeded != 0 { + t.Fatalf("summary.Succeeded = %d, want %d", summary.Succeeded, 0) + } + if summary.Failed != 1 { + t.Fatalf("summary.Failed = %d, want %d", summary.Failed, 1) + } + if len(summary.Results) != 1 { + t.Fatalf("len(summary.Results) = %d, want %d", len(summary.Results), 1) + } + if summary.Results[0].Status != storage.FetchRunStatusFailed { + t.Fatalf("summary.Results[0].Status = %q, want %q", summary.Results[0].Status, storage.FetchRunStatusFailed) + } + if !strings.Contains(summary.Results[0].ErrorMessage, "decode GitHub releases response") { + t.Fatalf("summary.Results[0].ErrorMessage = %q, want decode error", summary.Results[0].ErrorMessage) + } + + contentCount, err := queryCount(ctx, db, `SELECT COUNT(*) FROM content_items WHERE binding_id = ?`, "changelog-binding-github-1") + if err != nil { + t.Fatalf("queryCount(content_items) error = %v", err) + } + if contentCount != 0 { + t.Fatalf("contentCount = %d, want %d", contentCount, 0) + } +} + +func TestChangelogFetcherFailsOnMalformedRSS(t *testing.T) { + t.Parallel() + + ctx := context.Background() + discoveredAt := time.Date(2026, 4, 22, 16, 0, 0, 0, time.UTC) + + server := newChangelogHTTPTestServer(t, changelogHTTPTestServerOptions{ + Responses: map[string][]changelogHTTPResponse{ + "/feeds/changelog.xml": { + { + ContentType: "application/rss+xml; charset=utf-8", + Body: ``, + }, + }, + }, + }) + defer server.Close() + + pack := runtimeFixturePackWithChangelogBinding( + "changelog-rss-malformed", + "changelog-binding-rss-1", + config.JSONMap{ + "mode": "rss", + "feed_url": server.URL + "/feeds/changelog.xml", + }, + ) + + db := openRuntimeTestDB(t, pack) + defer db.Close() + + registry := NewRegistry() + registry.MustRegister(SourceKindChangelog, newChangelogTestFetcher(server)) + + runtime := Runtime{ + DB: db, + Registry: registry, + Now: sequenceClock( + discoveredAt, + discoveredAt.Add(time.Second), + ), + } + + summary, err := runtime.Execute(ctx, pack, Options{ + SourceKind: SourceKindChangelog, + BindingID: "changelog-binding-rss-1", + }) + if err != nil { + t.Fatalf("Execute() error = %v", err) + } + + if summary.Attempted != 1 { + t.Fatalf("summary.Attempted = %d, want %d", summary.Attempted, 1) + } + if summary.Succeeded != 0 { + t.Fatalf("summary.Succeeded = %d, want %d", summary.Succeeded, 0) + } + if summary.Failed != 1 { + t.Fatalf("summary.Failed = %d, want %d", summary.Failed, 1) + } + if len(summary.Results) != 1 { + t.Fatalf("len(summary.Results) = %d, want %d", len(summary.Results), 1) + } + if summary.Results[0].Status != storage.FetchRunStatusFailed { + t.Fatalf("summary.Results[0].Status = %q, want %q", summary.Results[0].Status, storage.FetchRunStatusFailed) + } + if !strings.Contains(summary.Results[0].ErrorMessage, "decode changelog RSS response") { + t.Fatalf("summary.Results[0].ErrorMessage = %q, want decode error", summary.Results[0].ErrorMessage) + } + + contentCount, err := queryCount(ctx, db, `SELECT COUNT(*) FROM content_items WHERE binding_id = ?`, "changelog-binding-rss-1") + if err != nil { + t.Fatalf("queryCount(content_items) error = %v", err) + } + if contentCount != 0 { + t.Fatalf("contentCount = %d, want %d", contentCount, 0) + } +} + +func newChangelogTestFetcher(server *httptest.Server) *ChangelogFetcher { + return NewChangelogFetcher(ChangelogFetcherOptions{ + GitHubBaseURL: server.URL + "/", + GitHubHTTPClient: server.Client(), + GitHubUserAgent: "signalscope-test", + RSSHTTPClient: server.Client(), + RSSUserAgent: "signalscope-test", + Now: func() time.Time { + return time.Date(2026, 4, 22, 0, 0, 0, 0, time.UTC) + }, + }) +} + +func runtimeFixturePackWithChangelogBinding(datasetID, bindingID string, scope config.JSONMap) config.Pack { + pack := runtimeFixturePack(datasetID) + + pack.Sources.Sources = append(pack.Sources.Sources, config.Source{ + ID: "changelog", + Kind: SourceKindChangelog, + Enabled: true, + Defaults: config.JSONMap{}, + }) + + pack.Sources.Bindings = append(pack.Sources.Bindings, config.Binding{ + ID: bindingID, + EntityID: "org-1", + SourceID: "changelog", + Enabled: true, + Scope: scope, + Notes: "changelog test binding", + }) + + return pack +} + +func readChangelogContentItemsForBinding(ctx context.Context, db *sql.DB, bindingID string) ([]changelogContentItemRow, error) { + rows, err := db.QueryContext( + ctx, + `SELECT + binding_id, + item_type, + title, + summary, + url, + external_id, + published_at, + discovered_at, + metadata_json + FROM content_items + WHERE binding_id = ? + ORDER BY id ASC`, + bindingID, + ) + if err != nil { + return nil, err + } + defer rows.Close() + + var result []changelogContentItemRow + for rows.Next() { + var row changelogContentItemRow + if err := rows.Scan( + &row.BindingID, + &row.ItemType, + &row.Title, + &row.Summary, + &row.URL, + &row.ExternalID, + &row.PublishedAt, + &row.DiscoveredAt, + &row.MetadataJSON, + ); err != nil { + return nil, err + } + + result = append(result, row) + } + + if err := rows.Err(); err != nil { + return nil, err + } + + return result, nil +} + +func newChangelogHTTPTestServer(t *testing.T, options changelogHTTPTestServerOptions) *httptest.Server { + t.Helper() + + responses := make(map[string][]changelogHTTPResponse, len(options.Responses)) + for route, sequence := range options.Responses { + responses[route] = append([]changelogHTTPResponse(nil), sequence...) + } + + var mu sync.Mutex + requestCounts := make(map[string]int, len(responses)) + + handler := http.HandlerFunc(func(w http.ResponseWriter, request *http.Request) { + route := request.URL.Path + if request.URL.RawQuery != "" { + route += "?" + request.URL.RawQuery + } + + mu.Lock() + sequence, ok := responses[route] + index := requestCounts[route] + if ok { + requestCounts[route]++ + } + mu.Unlock() + + if !ok { + http.Error(w, "unexpected route", http.StatusNotFound) + return + } + + if index >= len(sequence) { + index = len(sequence) - 1 + } + + response := sequence[index] + status := response.Status + if status == 0 { + status = http.StatusOK + } + + contentType := strings.TrimSpace(response.ContentType) + if contentType == "" { + trimmed := strings.TrimSpace(response.Body) + if strings.HasPrefix(trimmed, "<") { + contentType = "application/rss+xml; charset=utf-8" + } else { + contentType = "application/json; charset=utf-8" + } + } + + w.Header().Set("Content-Type", contentType) + w.WriteHeader(status) + _, _ = io.WriteString(w, response.Body) + }) + + return httptest.NewServer(handler) +} diff --git a/internal/fetch/runtime.go b/internal/fetch/runtime.go index 8ccfbef..86a302d 100644 --- a/internal/fetch/runtime.go +++ b/internal/fetch/runtime.go @@ -36,6 +36,7 @@ func DefaultRegistry() *Registry { registry := NewRegistry() registry.MustRegister(SourceKindGitHub, NewGitHubFetcher(GitHubFetcherOptions{})) registry.MustRegister(SourceKindNewsRSS, NewNewsRSSFetcher(NewsRSSFetcherOptions{})) + registry.MustRegister(SourceKindChangelog, NewChangelogFetcher(ChangelogFetcherOptions{})) return registry } diff --git a/internal/fetch/runtime_test.go b/internal/fetch/runtime_test.go index db544a9..7720ffc 100644 --- a/internal/fetch/runtime_test.go +++ b/internal/fetch/runtime_test.go @@ -187,7 +187,7 @@ func TestRuntimeExecuteMarksBindingsFailedWhenNoFetcherIsRegistered(t *testing.T t.Parallel() ctx := context.Background() - pack := runtimeFixturePackWithChangelog("runtime-fetch-no-fetcher") + pack := runtimeFixturePackWithSocial("runtime-fetch-no-fetcher") db := openRuntimeTestDB(t, pack) defer db.Close() @@ -200,7 +200,7 @@ func TestRuntimeExecuteMarksBindingsFailedWhenNoFetcherIsRegistered(t *testing.T ), } - summary, err := runtime.Execute(ctx, pack, Options{SourceKind: "changelog"}) + summary, err := runtime.Execute(ctx, pack, Options{SourceKind: "social"}) if err != nil { t.Fatalf("Execute() error = %v", err) } @@ -219,7 +219,7 @@ func TestRuntimeExecuteMarksBindingsFailedWhenNoFetcherIsRegistered(t *testing.T if result.Status != storage.FetchRunStatusFailed { t.Fatalf("result.Status = %q, want %q", result.Status, storage.FetchRunStatusFailed) } - if !strings.Contains(result.ErrorMessage, `no fetcher registered for source kind "changelog"`) { + if !strings.Contains(result.ErrorMessage, `no fetcher registered for source kind "social"`) { t.Fatalf("result.ErrorMessage = %q, want missing fetcher message", result.ErrorMessage) } } @@ -231,8 +231,8 @@ func TestRuntimeExecuteMarksBindingsFailedWhenNoFetcherIsRegistered(t *testing.T if len(rows) != 1 { t.Fatalf("len(rows) = %d, want %d", len(rows), 1) } - if rows[0].BindingID != "changelog-binding-1" { - t.Fatalf("rows[0].BindingID = %q, want %q", rows[0].BindingID, "changelog-binding-1") + if rows[0].BindingID != "social-binding-1" { + t.Fatalf("rows[0].BindingID = %q, want %q", rows[0].BindingID, "social-binding-1") } if rows[0].Status != storage.FetchRunStatusFailed { t.Fatalf("rows[0].Status = %q, want %q", rows[0].Status, storage.FetchRunStatusFailed) @@ -407,6 +407,30 @@ func runtimeFixturePackWithChangelog(datasetID string) config.Pack { return pack } +func runtimeFixturePackWithSocial(datasetID string) config.Pack { + pack := runtimeFixturePack(datasetID) + + pack.Sources.Sources = append(pack.Sources.Sources, config.Source{ + ID: "social", + Kind: "social", + Enabled: true, + Defaults: config.JSONMap{}, + }) + + pack.Sources.Bindings = append(pack.Sources.Bindings, config.Binding{ + ID: "social-binding-1", + EntityID: "org-1", + SourceID: "social", + Enabled: true, + Scope: config.JSONMap{ + "handle": "@example", + }, + Notes: "social binding 1", + }) + + return pack +} + func readFetchRuns(ctx context.Context, db *sql.DB) ([]fetchRunRow, error) { rows, err := db.QueryContext( ctx,