From 0d9c09618cf9c9f7f993feb9d9580f18a8cbf624 Mon Sep 17 00:00:00 2001 From: CriSTEM Date: Wed, 15 Apr 2026 20:19:57 -0600 Subject: [PATCH 1/2] chore(repo): normalize docs whitespace --- docs/configuration.md | 9 +++++++++ 1 file changed, 9 insertions(+) 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. From 3f8f27d1a917df09abbeedbbd000ff1eb486a654 Mon Sep 17 00:00:00 2001 From: CriSTEM Date: Wed, 15 Apr 2026 20:32:01 -0600 Subject: [PATCH 2/2] feat(fetch): ingest changelog GitHub releases and RSS feeds --- internal/fetch/changelog.go | 727 ++++++++++++++++++++++++++ internal/fetch/changelog_test.go | 839 +++++++++++++++++++++++++++++++ internal/fetch/runtime.go | 1 + internal/fetch/runtime_test.go | 34 +- 4 files changed, 1596 insertions(+), 5 deletions(-) create mode 100644 internal/fetch/changelog.go create mode 100644 internal/fetch/changelog_test.go 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,