Skip to content
Merged

Dev #30

Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 4 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,8 @@ own their Kick chat archive instead of depending on short-lived browser chat his
- Infinite-scroll message results with reply context, clickable links, and inline emotes.
- Public user and channel profile pages with activity summaries.
- Prediction pages that run client-side and do not store prediction data.
- Admin dashboard for followed channels, users, listener health, storage, and cleanup previews.
- Public request form for channel suggestions and feedback, reviewed from the admin dashboard.
- Admin dashboard for followed channels, requests, users, listener health, storage, and cleanup previews.
- Durable ingestion: raw events pass through NATS JetStream before ClickHouse normalization.
- Docker Compose runtime with Go, Next.js, NATS JetStream, ClickHouse, and SQLite.

Expand All @@ -76,7 +77,8 @@ The app is organized around:
message volume.
- **Prediction:** a public channel prediction view that fetches live Kick prediction data in the
browser without storing it.
- **Admin:** a protected dashboard for managing followed channels and watching ingestion health.
- **Requests:** a public form for channel suggestions and feedback.
- **Admin:** a protected dashboard for managing followed channels, requests, and ingestion health.

The default stack stores high-volume chat data in ClickHouse and keeps control-plane data such as
admins, followed channels, sender profiles, retention settings, and heartbeats in SQLite.
Expand Down
28 changes: 28 additions & 0 deletions apps/api-go/cmd/api/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import (
"net/http"
"os"
"os/signal"
"path/filepath"
"syscall"
"time"

Expand All @@ -22,15 +23,19 @@ import (
"github.com/YSelim0/kick-logs/apps/api-go/internal/infra/natsstream"
operationsinfra "github.com/YSelim0/kick-logs/apps/api-go/internal/infra/operations"
ratelimitinfra "github.com/YSelim0/kick-logs/apps/api-go/internal/infra/ratelimit"
"github.com/YSelim0/kick-logs/apps/api-go/internal/infra/snapshots"
sqliteinfra "github.com/YSelim0/kick-logs/apps/api-go/internal/infra/sqlite"
"github.com/YSelim0/kick-logs/apps/api-go/internal/ports"
analyticsusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/analytics"
authusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/auth"
channelsusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/channels"
datamanagementusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/data_management"
directoryusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/directory"
homepageusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/homepage"
kicksyncusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/kicksync"
messagesusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/messages"
profilesusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/profiles"
requestsusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/requests"
webhookprocessorusecase "github.com/YSelim0/kick-logs/apps/api-go/internal/usecase/webhookprocessor"
)

Expand Down Expand Up @@ -146,16 +151,21 @@ func main() {
}
var messageService *messagesusecase.Service
var analyticsService *analyticsusecase.Service
var homepageRepository ports.AnalyticsRepository
var profileService *profilesusecase.Service
var requestService *requestsusecase.Service
var subPeriodRepoForAPI ports.SubscriptionPeriodRepository
if clickHouseConn != nil {
messageRepository := clickhouseinfra.NewMessageRepository(clickHouseConn)
analyticsRepository := clickhouseinfra.NewAnalyticsRepository(clickHouseConn)
subPeriodRepo := clickhouseinfra.NewSubscriptionPeriodRepository(clickHouseConn)
userRequestRepo := clickhouseinfra.NewUserRequestRepository(clickHouseConn)
subPeriodRepoForAPI = subPeriodRepo
messageService = messagesusecase.NewService(messageRepository)
analyticsService = analyticsusecase.NewService(analyticsRepository)
homepageRepository = clickhouseinfra.NewHomepageAnalyticsRepository(clickHouseConn)
profileService = profilesusecase.NewService(analyticsRepository, channelRepo, senderRepo)
requestService = requestsusecase.NewService(userRequestRepo)

processorSvc := webhookprocessorusecase.NewService(
logger,
Expand All @@ -167,6 +177,21 @@ func main() {
)
processorSvc.Start(context.Background())
}
homepageService := homepageusecase.NewService(
homepageRepository,
snapshots.NewHomepageStore(filepath.Join(filepath.Dir(cfg.SQLitePath), "homepage-analytics-v1.json")),
logger,
)
homepageCtx, cancelHomepage := context.WithCancel(context.Background())
homepageDone := make(chan struct{})
go func() {
defer close(homepageDone)
homepageService.Run(homepageCtx)
}()
defer func() {
cancelHomepage()
<-homepageDone
}()
operationsRepo := operationsinfra.NewRepository(
sqliteDB,
cfg.SQLitePath,
Expand All @@ -181,11 +206,14 @@ func main() {
Config: cfg,
Auth: authService,
Analytics: analyticsService,
Homepage: homepageService,
Channels: channelService,
Messages: messageService,
Profiles: profileService,
Data: dataManagementService,
Directory: directoryusecase.NewService(sqliteinfra.NewDirectoryRepository(sqliteDB)),
KickSync: kickSyncService,
Requests: requestService,
WebhookEvents: webhookEventRepo,
WebhookVerifier: webhookVerifier,
WebhookEventSubs: eventSubRepo,
Expand Down
32 changes: 32 additions & 0 deletions apps/api-go/internal/domain/directory.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
package domain

type DirectoryKind string

const (
DirectoryUsers DirectoryKind = "users"
DirectoryChannels DirectoryKind = "channels"
)

type DirectoryIdentity struct {
ID int64
Name string
Slug string
ProfileImageURL string
SortSlug string
}

type DirectoryPosition struct {
Slug string
ID int64
}

type DirectoryQuery struct {
Prefix string
Limit int
After DirectoryPosition
}

type DirectoryPage struct {
Items []DirectoryIdentity
NextCursor string
}
16 changes: 16 additions & 0 deletions apps/api-go/internal/domain/homepage.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
package domain

import "time"

// HomepageSnapshot is a disposable read model, never the source of message history.
type HomepageSnapshot struct {
Version int
AsOf time.Time
Start time.Time
End time.Time
Overview AnalyticsOverview
Volume []MessageVolumePoint
TopChannels []TopChannelAnalytics
TopSenders []TopSenderAnalytics
TopEmotes []TopEmoteAnalytics
}
100 changes: 100 additions & 0 deletions apps/api-go/internal/domain/models.go
Original file line number Diff line number Diff line change
Expand Up @@ -549,3 +549,103 @@ type ChannelSubscriptionSummary struct {
ActiveGiftedCount int64
LatestEventAt time.Time
}

type ChannelSubscriberFilter struct {
FollowedChannelID int64
GiftOnly bool
Limit uint64
Offset uint64
}

type ChannelSubscriber struct {
SubscriberKickUserID int64
Username string
Slug string
ProfileImageURL string
IsGift bool
GifterKickUserID int64
GifterUsername string
GifterSlug string
GifterProfileImageURL string
StartedAt time.Time
ExpiresAt time.Time
}

type ChannelSubscriberPage struct {
Items []ChannelSubscriber
Count int64
Limit uint64
Offset uint64
}

type UserRequestType string

const (
UserRequestTypeChannelRequest UserRequestType = "channel_request"
UserRequestTypeFeedback UserRequestType = "feedback"
)

type UserRequestStatus string

const (
UserRequestStatusNew UserRequestStatus = "new"
UserRequestStatusReviewing UserRequestStatus = "reviewing"
UserRequestStatusApproved UserRequestStatus = "approved"
UserRequestStatusRejected UserRequestStatus = "rejected"
UserRequestStatusDone UserRequestStatus = "done"
UserRequestStatusDuplicate UserRequestStatus = "duplicate"
)

type UserRequestEventType string

const (
UserRequestEventStatusChanged UserRequestEventType = "status_changed"
UserRequestEventNoteAdded UserRequestEventType = "note_added"
UserRequestEventArchived UserRequestEventType = "archived"
)

type UserRequest struct {
ID string
Type UserRequestType
Title string
Message string
ChannelSlug string
ChannelDisplayName string
Contact string
IPHash string
UserAgentHash string
CreatedAt time.Time
}

type UserRequestEvent struct {
ID string
RequestID string
EventType UserRequestEventType
Status UserRequestStatus
Note string
AdminID int64
CreatedAt time.Time
}

type UserRequestState struct {
Request UserRequest
CurrentStatus UserRequestStatus
IsArchived bool
LatestEventAt time.Time
}

type UserRequestDetail struct {
State UserRequestState
Events []UserRequestEvent
}

type UserRequestListFilter struct {
Type UserRequestType
Status UserRequestStatus
Archived *bool
Query string
Start time.Time
End time.Time
Limit uint64
Offset uint64
}
35 changes: 34 additions & 1 deletion apps/api-go/internal/http/analytics_profile_routes_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package httpapi
import (
"context"
"database/sql"
"encoding/json"
"errors"
"io"
"log/slog"
Expand Down Expand Up @@ -130,6 +131,36 @@ func TestProfileRoutesCacheChannelAnalyticsBySlug(t *testing.T) {
}
}

func TestProfileRoutesReturnFourteenDayVolumeSeries(t *testing.T) {
handler, analyticsRepo := newAnalyticsProfileTestRouterWithRepo()

response := httptest.NewRecorder()
handler.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/channels/hype/analytics", nil))
if response.Code != http.StatusOK {
t.Fatalf("status = %d body = %s", response.Code, response.Body.String())
}

var payload struct {
MessageVolume []struct {
BucketStart string `json:"bucket_start"`
MessageCount int64 `json:"message_count"`
} `json:"message_volume"`
}
if err := json.Unmarshal(response.Body.Bytes(), &payload); err != nil {
t.Fatalf("decode response: %v", err)
}
if len(payload.MessageVolume) != 14 {
t.Fatalf("message volume length = %d body = %s", len(payload.MessageVolume), response.Body.String())
}
if len(analyticsRepo.volumeFilters) == 0 {
t.Fatal("expected profile route to request message volume")
}
filter := analyticsRepo.volumeFilters[len(analyticsRepo.volumeFilters)-1]
if filter.Start.IsZero() || filter.End.IsZero() {
t.Fatalf("volume filter = %#v", filter)
}
}

func newAnalyticsProfileTestRouter() http.Handler {
handler, _ := newAnalyticsProfileTestRouterWithRepo()
return handler
Expand All @@ -152,6 +183,7 @@ type fakeAnalyticsRepository struct {
now time.Time
topEmotesErr error
topEmotesCalls int
volumeFilters []domain.AnalyticsFilter
}

func newFakeAnalyticsRepository() *fakeAnalyticsRepository {
Expand All @@ -174,9 +206,10 @@ func (repo *fakeAnalyticsRepository) Overview(

func (repo *fakeAnalyticsRepository) MessageVolume(
_ context.Context,
_ domain.AnalyticsFilter,
filter domain.AnalyticsFilter,
_ domain.AnalyticsBucket,
) ([]domain.MessageVolumePoint, error) {
repo.volumeFilters = append(repo.volumeFilters, filter)
return []domain.MessageVolumePoint{{BucketStart: repo.now, MessageCount: 2}}, nil
}

Expand Down
Loading
Loading