feat: connect explore datasets to backend

This commit is contained in:
2026-07-23 15:07:35 +08:00
parent 0ca2898ac0
commit 26dd3004a1
14 changed files with 1445 additions and 196 deletions

View File

@@ -0,0 +1,353 @@
package dataset
import (
"errors"
"net/http"
"strings"
"time"
"github.com/gin-gonic/gin"
"gorm.io/gorm"
"senlinai-agent/backend/internal/httpx"
"senlinai-agent/backend/internal/logic/auth"
"senlinai-agent/backend/internal/models"
)
type Handler struct {
service *Service
}
type SourceDTO struct {
ID string `json:"id"`
Name string `json:"name"`
Kind string `json:"kind"`
URL string `json:"url"`
Description string `json:"description"`
Enabled bool `json:"enabled"`
LastSyncedAt *time.Time `json:"lastSyncedAt"`
ItemCount int64 `json:"itemCount"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}
type ItemDTO struct {
ID string `json:"id"`
SourceID string `json:"sourceId"`
Title string `json:"title"`
Summary string `json:"summary"`
Content string `json:"content"`
URL string `json:"url"`
Status string `json:"status"`
Starred bool `json:"starred"`
PublishedAt *time.Time `json:"publishedAt"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}
type CronDTO struct {
ID string `json:"id"`
SourceID string `json:"sourceId"`
Schedule string `json:"schedule"`
Status string `json:"status"`
Enabled bool `json:"enabled"`
NextRunAt *time.Time `json:"nextRunAt"`
LastRunAt *time.Time `json:"lastRunAt"`
LastResult string `json:"lastResult"`
CreatedAt time.Time `json:"createdAt"`
UpdatedAt time.Time `json:"updatedAt"`
}
type sourceRequest struct {
Name string `json:"name"`
Kind string `json:"kind"`
URL string `json:"url"`
Description string `json:"description"`
Enabled *bool `json:"enabled"`
}
func NewHandler(service *Service) *Handler {
return &Handler{service: service}
}
func (h *Handler) Register(router gin.IRouter) {
router.GET("/dataset-sources", h.listSources)
router.POST("/dataset-sources", h.createSource)
router.PATCH("/dataset-sources/:id", h.updateSource)
router.DELETE("/dataset-sources/:id", h.deleteSource)
router.GET("/dataset-items", h.listItems)
router.POST("/dataset-items", h.createItem)
router.PATCH("/dataset-items/:id", h.updateItem)
router.GET("/dataset-crons", h.listCrons)
router.POST("/dataset-crons/sync", h.queueSync)
}
func (h *Handler) listSources(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
records, err := h.service.ListSources(userID)
if err != nil {
writeError(c, err)
return
}
items := make([]SourceDTO, 0, len(records))
for _, record := range records {
items = append(items, sourceDTO(record.Source, record.ItemCount))
}
c.JSON(http.StatusOK, items)
}
func (h *Handler) createSource(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
var input sourceRequest
if err := c.ShouldBindJSON(&input); err != nil {
writeInvalidRequest(c)
return
}
enabled := true
if input.Enabled != nil {
enabled = *input.Enabled
}
source, err := h.service.CreateSource(userID, SourceInput{
Name: input.Name, Kind: input.Kind, URL: input.URL, Description: input.Description, Enabled: enabled,
})
if err != nil {
writeError(c, err)
return
}
c.JSON(http.StatusCreated, sourceDTO(*source, 0))
}
func (h *Handler) updateSource(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
identity, ok := httpx.IdentityParam(c, "id")
if !ok {
return
}
var input sourceRequest
if err := c.ShouldBindJSON(&input); err != nil {
writeInvalidRequest(c)
return
}
enabled := true
if input.Enabled != nil {
enabled = *input.Enabled
}
source, err := h.service.UpdateSource(userID, identity, SourceInput{
Name: input.Name, Kind: input.Kind, URL: input.URL, Description: input.Description, Enabled: enabled,
})
if err != nil {
writeError(c, err)
return
}
itemCount, err := h.service.CountItems(userID, source.ID)
if err != nil {
writeError(c, err)
return
}
c.JSON(http.StatusOK, sourceDTO(*source, itemCount))
}
func (h *Handler) deleteSource(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
identity, ok := httpx.IdentityParam(c, "id")
if !ok {
return
}
if err := h.service.DeleteSource(userID, identity); err != nil {
writeError(c, err)
return
}
c.Status(http.StatusNoContent)
}
func (h *Handler) listItems(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
items, err := h.service.ListItems(userID, c.Query("sourceId"))
if err != nil {
writeError(c, err)
return
}
result := make([]ItemDTO, 0, len(items))
for _, item := range items {
result = append(result, itemDTO(item))
}
c.JSON(http.StatusOK, result)
}
func (h *Handler) createItem(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
var input struct {
SourceID string `json:"sourceId"`
Title string `json:"title"`
Summary string `json:"summary"`
Content string `json:"content"`
URL string `json:"url"`
PublishedAt string `json:"publishedAt"`
}
if err := c.ShouldBindJSON(&input); err != nil {
writeInvalidRequest(c)
return
}
publishedAt, err := parseOptionalTime(input.PublishedAt)
if err != nil {
writeInvalidRequest(c)
return
}
item, err := h.service.CreateItem(userID, ItemInput{
SourceIdentity: input.SourceID, Title: input.Title, Summary: input.Summary,
Content: input.Content, URL: input.URL, PublishedAt: publishedAt,
})
if err != nil {
writeError(c, err)
return
}
c.JSON(http.StatusCreated, itemDTO(*item))
}
func (h *Handler) updateItem(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
identity, ok := httpx.IdentityParam(c, "id")
if !ok {
return
}
var input struct {
Status string `json:"status"`
Starred bool `json:"starred"`
}
if err := c.ShouldBindJSON(&input); err != nil {
writeInvalidRequest(c)
return
}
item, err := h.service.UpdateItem(userID, identity, ItemUpdate{Status: input.Status, Starred: input.Starred})
if err != nil {
writeError(c, err)
return
}
c.JSON(http.StatusOK, itemDTO(*item))
}
func (h *Handler) listCrons(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
crons, err := h.service.ListCrons(userID)
if err != nil {
writeError(c, err)
return
}
result := make([]CronDTO, 0, len(crons))
for _, cron := range crons {
result = append(result, cronDTO(cron))
}
c.JSON(http.StatusOK, result)
}
func (h *Handler) queueSync(c *gin.Context) {
userID, ok := currentUser(c)
if !ok {
return
}
crons, err := h.service.QueueSync(userID)
if err != nil {
writeError(c, err)
return
}
result := make([]CronDTO, 0, len(crons))
for _, cron := range crons {
result = append(result, cronDTO(cron))
}
c.JSON(http.StatusAccepted, result)
}
func currentUser(c *gin.Context) (uint, bool) {
userID, ok := auth.CurrentUserID(c)
if !ok {
httpx.Error(c, http.StatusUnauthorized, "unauthorized", "未登录或登录已失效")
return 0, false
}
return userID, true
}
func writeError(c *gin.Context, err error) {
switch {
case errors.Is(err, gorm.ErrRecordNotFound):
httpx.Error(c, http.StatusNotFound, "not_found", "数据源或数据条目不存在")
case errors.Is(err, ErrNameRequired), errors.Is(err, ErrKindInvalid),
errors.Is(err, ErrURLInvalid), errors.Is(err, ErrTitleRequired), errors.Is(err, ErrStatusInvalid):
writeInvalidRequest(c)
default:
httpx.Error(c, http.StatusInternalServerError, "internal_error", "探索数据操作失败")
}
}
func writeInvalidRequest(c *gin.Context) {
httpx.Error(c, http.StatusBadRequest, "invalid_request", "请求参数无效")
}
func sourceDTO(source models.SaDatasetSource, itemCount int64) SourceDTO {
return SourceDTO{
ID: source.Identity, Name: source.Name, Kind: source.Kind, URL: source.URL,
Description: source.Description, Enabled: source.Enabled,
LastSyncedAt: utcTime(source.LastSyncedAt), ItemCount: itemCount,
CreatedAt: source.CreatedAt.UTC(), UpdatedAt: source.UpdatedAt.UTC(),
}
}
func itemDTO(item models.SaDatasetItem) ItemDTO {
return ItemDTO{
ID: item.Identity, SourceID: item.SourceIdentity, Title: item.Title,
Summary: item.Summary, Content: item.Content, URL: item.URL, Status: item.Status,
Starred: item.Starred, PublishedAt: utcTime(item.PublishedAt),
CreatedAt: item.CreatedAt.UTC(), UpdatedAt: item.UpdatedAt.UTC(),
}
}
func cronDTO(cron models.SaDatasetCron) CronDTO {
return CronDTO{
ID: cron.Identity, SourceID: cron.SourceIdentity, Schedule: cron.Schedule,
Status: cron.Status, Enabled: cron.Enabled, NextRunAt: utcTime(cron.NextRunAt),
LastRunAt: utcTime(cron.LastRunAt), LastResult: cron.LastResult,
CreatedAt: cron.CreatedAt.UTC(), UpdatedAt: cron.UpdatedAt.UTC(),
}
}
func parseOptionalTime(value string) (*time.Time, error) {
if strings.TrimSpace(value) == "" {
return nil, nil
}
parsed, err := time.Parse(time.RFC3339, value)
if err != nil {
return nil, err
}
parsed = parsed.UTC()
return &parsed, nil
}
func utcTime(value *time.Time) *time.Time {
if value == nil {
return nil
}
result := value.UTC()
return &result
}

View File

@@ -0,0 +1,139 @@
package dataset
import (
"bytes"
"encoding/json"
"fmt"
"net/http"
"net/http/httptest"
"testing"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/require"
"gorm.io/gorm"
"senlinai-agent/backend/internal/config"
"senlinai-agent/backend/internal/httpx"
"senlinai-agent/backend/internal/models"
)
func TestDatasetAPIConnectsSourcesItemsAndSyncTasks(t *testing.T) {
database := newDatasetTestDatabase(t)
owner := createDatasetTestUser(t, database, "owner@example.com")
other := createDatasetTestUser(t, database, "other@example.com")
ownerRouter := datasetTestRouter(database, owner.ID)
invalid := performDatasetRequest(t, ownerRouter, http.MethodPost, "/api/v1/dataset-sources", map[string]any{
"name": "Invalid RSS", "kind": "rss", "url": "",
})
require.Equal(t, http.StatusBadRequest, invalid.Code)
createdSource := performDatasetRequest(t, ownerRouter, http.MethodPost, "/api/v1/dataset-sources", map[string]any{
"name": "Industry feed", "kind": "rss", "url": "https://example.com/feed",
"description": "Industry updates", "enabled": true,
})
require.Equal(t, http.StatusCreated, createdSource.Code)
var source SourceDTO
require.NoError(t, json.Unmarshal(createdSource.Body.Bytes(), &source))
require.NotEmpty(t, source.ID)
require.Equal(t, "Industry feed", source.Name)
createdItem := performDatasetRequest(t, ownerRouter, http.MethodPost, "/api/v1/dataset-items", map[string]any{
"sourceId": source.ID, "title": "A collected article", "summary": "Summary",
"url": "https://example.com/article",
})
require.Equal(t, http.StatusCreated, createdItem.Code)
var item ItemDTO
require.NoError(t, json.Unmarshal(createdItem.Body.Bytes(), &item))
require.Equal(t, source.ID, item.SourceID)
listSources := performDatasetRequest(t, ownerRouter, http.MethodGet, "/api/v1/dataset-sources", nil)
require.Equal(t, http.StatusOK, listSources.Code)
var sources []SourceDTO
require.NoError(t, json.Unmarshal(listSources.Body.Bytes(), &sources))
require.Len(t, sources, 1)
require.Equal(t, int64(1), sources[0].ItemCount)
updatedItem := performDatasetRequest(t, ownerRouter, http.MethodPatch, "/api/v1/dataset-items/"+item.ID, map[string]any{
"status": "read", "starred": true,
})
require.Equal(t, http.StatusOK, updatedItem.Code)
require.NoError(t, json.Unmarshal(updatedItem.Body.Bytes(), &item))
require.Equal(t, "read", item.Status)
require.True(t, item.Starred)
firstSync := performDatasetRequest(t, ownerRouter, http.MethodPost, "/api/v1/dataset-crons/sync", map[string]any{})
require.Equal(t, http.StatusAccepted, firstSync.Code)
var firstCrons []CronDTO
require.NoError(t, json.Unmarshal(firstSync.Body.Bytes(), &firstCrons))
require.Len(t, firstCrons, 1)
require.Equal(t, "pending", firstCrons[0].Status)
secondSync := performDatasetRequest(t, ownerRouter, http.MethodPost, "/api/v1/dataset-crons/sync", map[string]any{})
require.Equal(t, http.StatusAccepted, secondSync.Code)
var secondCrons []CronDTO
require.NoError(t, json.Unmarshal(secondSync.Body.Bytes(), &secondCrons))
require.Equal(t, firstCrons[0].ID, secondCrons[0].ID)
otherRouter := datasetTestRouter(database, other.ID)
otherSources := performDatasetRequest(t, otherRouter, http.MethodGet, "/api/v1/dataset-sources", nil)
require.Equal(t, http.StatusOK, otherSources.Code)
require.JSONEq(t, `[]`, otherSources.Body.String())
otherUpdate := performDatasetRequest(t, otherRouter, http.MethodPatch, "/api/v1/dataset-sources/"+source.ID, map[string]any{
"name": "Stolen", "kind": "manual", "enabled": true,
})
require.Equal(t, http.StatusNotFound, otherUpdate.Code)
deleted := performDatasetRequest(t, ownerRouter, http.MethodDelete, "/api/v1/dataset-sources/"+source.ID, nil)
require.Equal(t, http.StatusNoContent, deleted.Code)
var itemCount, cronCount int64
require.NoError(t, database.Model(&models.SaDatasetItem{}).Count(&itemCount).Error)
require.NoError(t, database.Model(&models.SaDatasetCron{}).Count(&cronCount).Error)
require.Zero(t, itemCount)
require.Zero(t, cronCount)
}
func newDatasetTestDatabase(t *testing.T) *gorm.DB {
t.Helper()
database, err := gorm.Open(sqlite.Open(fmt.Sprintf("file:%s?mode=memory&cache=shared", t.Name())), &gorm.Config{})
require.NoError(t, err)
require.NoError(t, database.AutoMigrate(
&models.SaUser{},
&models.SaDatasetSource{},
&models.SaDatasetItem{},
&models.SaDatasetCron{},
))
return database
}
func createDatasetTestUser(t *testing.T, database *gorm.DB, email string) models.SaUser {
t.Helper()
user := models.SaUser{Email: email, DisplayName: email, PasswordHash: "hash", Role: "user"}
require.NoError(t, database.Create(&user).Error)
return user
}
func datasetTestRouter(database *gorm.DB, userID uint) http.Handler {
return httpx.NewProtectedRouter(
config.Config{Env: "test"},
func(string) (uint, error) { return userID, nil },
NewHandler(NewService(database)),
)
}
func performDatasetRequest(t *testing.T, router http.Handler, method, path string, body any) *httptest.ResponseRecorder {
t.Helper()
var encoded []byte
if body != nil {
var err error
encoded, err = json.Marshal(body)
require.NoError(t, err)
}
request := httptest.NewRequest(method, path, bytes.NewReader(encoded))
request.Header.Set("Authorization", "Bearer test-token")
if body != nil {
request.Header.Set("Content-Type", "application/json")
}
recorder := httptest.NewRecorder()
router.ServeHTTP(recorder, request)
return recorder
}

View File

@@ -0,0 +1,278 @@
package dataset
import (
"errors"
"net/url"
"strings"
"time"
"gorm.io/gorm"
"senlinai-agent/backend/internal/models"
)
var (
ErrNameRequired = errors.New("dataset source name is required")
ErrKindInvalid = errors.New("dataset source kind is invalid")
ErrURLInvalid = errors.New("dataset source url is invalid")
ErrTitleRequired = errors.New("dataset item title is required")
ErrStatusInvalid = errors.New("dataset item status is invalid")
)
type Service struct {
db *gorm.DB
}
type SourceInput struct {
Name string
Kind string
URL string
Description string
Enabled bool
}
type ItemInput struct {
SourceIdentity string
Title string
Summary string
Content string
URL string
PublishedAt *time.Time
}
type ItemUpdate struct {
Status string
Starred bool
}
type SourceRecord struct {
Source models.SaDatasetSource
ItemCount int64
}
func NewService(database *gorm.DB) *Service {
return &Service{db: database}
}
func (s *Service) ListSources(userID uint) ([]SourceRecord, error) {
var sources []models.SaDatasetSource
if err := s.db.Where("created_by = ?", userID).Order("created_at asc, id asc").Find(&sources).Error; err != nil {
return nil, err
}
counts := make(map[uint]int64)
if len(sources) > 0 {
sourceIDs := make([]uint, 0, len(sources))
for _, source := range sources {
sourceIDs = append(sourceIDs, source.ID)
}
var rows []struct {
SourceID uint
Count int64
}
if err := s.db.Model(&models.SaDatasetItem{}).
Select("source_id, COUNT(*) AS count").
Where("created_by = ? AND source_id IN ?", userID, sourceIDs).
Group("source_id").
Scan(&rows).Error; err != nil {
return nil, err
}
for _, row := range rows {
counts[row.SourceID] = row.Count
}
}
result := make([]SourceRecord, 0, len(sources))
for _, source := range sources {
result = append(result, SourceRecord{Source: source, ItemCount: counts[source.ID]})
}
return result, nil
}
func (s *Service) CreateSource(userID uint, input SourceInput) (*models.SaDatasetSource, error) {
normalized, err := normalizeSourceInput(input)
if err != nil {
return nil, err
}
source := &models.SaDatasetSource{
CreatedBy: userID, Name: normalized.Name, Kind: normalized.Kind,
URL: normalized.URL, Description: normalized.Description, Enabled: normalized.Enabled,
}
return source, s.db.Create(source).Error
}
func (s *Service) UpdateSource(userID uint, identity string, input SourceInput) (*models.SaDatasetSource, error) {
normalized, err := normalizeSourceInput(input)
if err != nil {
return nil, err
}
var source models.SaDatasetSource
if err := s.db.Where("identity = ? AND created_by = ?", identity, userID).First(&source).Error; err != nil {
return nil, err
}
if err := s.db.Model(&source).Updates(map[string]any{
"name": normalized.Name, "kind": normalized.Kind, "url": normalized.URL,
"description": normalized.Description, "enabled": normalized.Enabled,
}).Error; err != nil {
return nil, err
}
return &source, nil
}
func (s *Service) DeleteSource(userID uint, identity string) error {
return s.db.Transaction(func(tx *gorm.DB) error {
var source models.SaDatasetSource
if err := tx.Where("identity = ? AND created_by = ?", identity, userID).First(&source).Error; err != nil {
return err
}
if err := tx.Where("source_id = ? AND created_by = ?", source.ID, userID).Delete(&models.SaDatasetCron{}).Error; err != nil {
return err
}
if err := tx.Where("source_id = ? AND created_by = ?", source.ID, userID).Delete(&models.SaDatasetItem{}).Error; err != nil {
return err
}
return tx.Delete(&source).Error
})
}
func (s *Service) ListItems(userID uint, sourceIdentity string) ([]models.SaDatasetItem, error) {
query := s.db.Where("created_by = ?", userID)
if sourceIdentity = strings.TrimSpace(sourceIdentity); sourceIdentity != "" {
source, err := s.findOwnedSource(userID, sourceIdentity)
if err != nil {
return nil, err
}
query = query.Where("source_id = ?", source.ID)
}
var items []models.SaDatasetItem
if err := query.Order("COALESCE(published_at, created_at) desc, id desc").Limit(200).Find(&items).Error; err != nil {
return nil, err
}
if items == nil {
items = []models.SaDatasetItem{}
}
return items, nil
}
func (s *Service) CountItems(userID uint, sourceID uint) (int64, error) {
var count int64
err := s.db.Model(&models.SaDatasetItem{}).
Where("created_by = ? AND source_id = ?", userID, sourceID).
Count(&count).Error
return count, err
}
func (s *Service) CreateItem(userID uint, input ItemInput) (*models.SaDatasetItem, error) {
title := strings.TrimSpace(input.Title)
if title == "" {
return nil, ErrTitleRequired
}
source, err := s.findOwnedSource(userID, input.SourceIdentity)
if err != nil {
return nil, err
}
itemURL := strings.TrimSpace(input.URL)
if itemURL != "" && !validHTTPURL(itemURL) {
return nil, ErrURLInvalid
}
item := &models.SaDatasetItem{
SourceID: source.ID, CreatedBy: userID, Title: title,
Summary: strings.TrimSpace(input.Summary), Content: strings.TrimSpace(input.Content),
URL: itemURL, Status: "unread", PublishedAt: utcOptionalTime(input.PublishedAt),
}
return item, s.db.Create(item).Error
}
func (s *Service) UpdateItem(userID uint, identity string, input ItemUpdate) (*models.SaDatasetItem, error) {
status := strings.TrimSpace(input.Status)
if status != "unread" && status != "read" && status != "archived" {
return nil, ErrStatusInvalid
}
var item models.SaDatasetItem
if err := s.db.Where("identity = ? AND created_by = ?", identity, userID).First(&item).Error; err != nil {
return nil, err
}
if err := s.db.Model(&item).Updates(map[string]any{"status": status, "starred": input.Starred}).Error; err != nil {
return nil, err
}
return &item, nil
}
func (s *Service) ListCrons(userID uint) ([]models.SaDatasetCron, error) {
var crons []models.SaDatasetCron
if err := s.db.Where("created_by = ?", userID).Order("created_at desc, id desc").Limit(100).Find(&crons).Error; err != nil {
return nil, err
}
if crons == nil {
crons = []models.SaDatasetCron{}
}
return crons, nil
}
func (s *Service) QueueSync(userID uint) ([]models.SaDatasetCron, error) {
result := make([]models.SaDatasetCron, 0)
err := s.db.Transaction(func(tx *gorm.DB) error {
var sources []models.SaDatasetSource
if err := tx.Where("created_by = ? AND enabled = ?", userID, true).Order("id asc").Find(&sources).Error; err != nil {
return err
}
now := time.Now().UTC()
for _, source := range sources {
var cron models.SaDatasetCron
err := tx.Where("source_id = ? AND created_by = ? AND status = ?", source.ID, userID, "pending").First(&cron).Error
if err == nil {
result = append(result, cron)
continue
}
if !errors.Is(err, gorm.ErrRecordNotFound) {
return err
}
cron = models.SaDatasetCron{
SourceID: source.ID, CreatedBy: userID, Schedule: "@once",
Status: "pending", Enabled: true, NextRunAt: &now,
}
if err := tx.Create(&cron).Error; err != nil {
return err
}
result = append(result, cron)
}
return nil
})
return result, err
}
func (s *Service) findOwnedSource(userID uint, identity string) (*models.SaDatasetSource, error) {
var source models.SaDatasetSource
err := s.db.Where("identity = ? AND created_by = ?", strings.TrimSpace(identity), userID).First(&source).Error
return &source, err
}
func normalizeSourceInput(input SourceInput) (SourceInput, error) {
input.Name = strings.TrimSpace(input.Name)
input.Kind = strings.ToLower(strings.TrimSpace(input.Kind))
input.URL = strings.TrimSpace(input.URL)
input.Description = strings.TrimSpace(input.Description)
if input.Name == "" {
return input, ErrNameRequired
}
if input.Kind != "manual" && input.Kind != "link" && input.Kind != "rss" {
return input, ErrKindInvalid
}
if input.URL != "" && !validHTTPURL(input.URL) {
return input, ErrURLInvalid
}
if input.Kind != "manual" && input.URL == "" {
return input, ErrURLInvalid
}
return input, nil
}
func validHTTPURL(value string) bool {
parsed, err := url.ParseRequestURI(value)
return err == nil && (parsed.Scheme == "http" || parsed.Scheme == "https") && parsed.Host != ""
}
func utcOptionalTime(value *time.Time) *time.Time {
if value == nil {
return nil
}
result := value.UTC()
return &result
}

View File

@@ -0,0 +1,24 @@
package models
import "time"
type SaDatasetCron struct {
ID uint `gorm:"primaryKey"`
Identity string `gorm:"type:char(36);uniqueIndex"`
SourceID uint `gorm:"index;not null"`
SourceIdentity string `gorm:"type:char(36);index"`
CreatedBy uint `gorm:"index;not null"`
CreatedByIdentity string `gorm:"type:char(36);index"`
Schedule string `gorm:"size:100;not null"`
Status string `gorm:"size:32;not null;default:pending;index"`
Enabled bool `gorm:"not null;default:true;index"`
NextRunAt *time.Time
LastRunAt *time.Time
LastResult string `gorm:"type:text"`
CreatedAt time.Time
UpdatedAt time.Time
}
func (SaDatasetCron) TableName() string {
return "sa_dataset_crons"
}

View File

@@ -0,0 +1,25 @@
package models
import "time"
type SaDatasetItem struct {
ID uint `gorm:"primaryKey"`
Identity string `gorm:"type:char(36);uniqueIndex"`
SourceID uint `gorm:"index;not null"`
SourceIdentity string `gorm:"type:char(36);index"`
CreatedBy uint `gorm:"index;not null"`
CreatedByIdentity string `gorm:"type:char(36);index"`
Title string `gorm:"size:500;not null"`
Summary string `gorm:"type:text"`
Content string `gorm:"type:text"`
URL string `gorm:"size:2048"`
Status string `gorm:"size:32;not null;default:unread;index"`
Starred bool `gorm:"not null;default:false;index"`
PublishedAt *time.Time
CreatedAt time.Time
UpdatedAt time.Time
}
func (SaDatasetItem) TableName() string {
return "sa_dataset_items"
}

View File

@@ -0,0 +1,22 @@
package models
import "time"
type SaDatasetSource struct {
ID uint `gorm:"primaryKey"`
Identity string `gorm:"type:char(36);uniqueIndex"`
CreatedBy uint `gorm:"index;not null"`
CreatedByIdentity string `gorm:"type:char(36);index"`
Name string `gorm:"size:160;not null"`
Kind string `gorm:"size:32;not null"`
URL string `gorm:"size:2048"`
Description string `gorm:"type:text"`
Enabled bool `gorm:"not null;default:true;index"`
LastSyncedAt *time.Time
CreatedAt time.Time
UpdatedAt time.Time
}
func (SaDatasetSource) TableName() string {
return "sa_dataset_sources"
}

View File

@@ -18,6 +18,33 @@ func (m *SaProject) BeforeCreate(tx *gorm.DB) error {
return resolveIdentity(tx, &SaUser{}, m.OwnerID, &m.OwnerIdentity)
}
func (m *SaDatasetSource) BeforeCreate(tx *gorm.DB) error {
if err := ensureIdentity(&m.Identity); err != nil {
return err
}
return resolveIdentity(tx, &SaUser{}, m.CreatedBy, &m.CreatedByIdentity)
}
func (m *SaDatasetItem) BeforeCreate(tx *gorm.DB) error {
if err := ensureIdentity(&m.Identity); err != nil {
return err
}
if err := resolveIdentity(tx, &SaDatasetSource{}, m.SourceID, &m.SourceIdentity); err != nil {
return err
}
return resolveIdentity(tx, &SaUser{}, m.CreatedBy, &m.CreatedByIdentity)
}
func (m *SaDatasetCron) BeforeCreate(tx *gorm.DB) error {
if err := ensureIdentity(&m.Identity); err != nil {
return err
}
if err := resolveIdentity(tx, &SaDatasetSource{}, m.SourceID, &m.SourceIdentity); err != nil {
return err
}
return resolveIdentity(tx, &SaUser{}, m.CreatedBy, &m.CreatedByIdentity)
}
func (m *SaInboxItem) BeforeCreate(tx *gorm.DB) error {
if err := ensureIdentity(&m.Identity); err != nil {
return err

View File

@@ -31,6 +31,9 @@ func AutoMigrate(database *gorm.DB) error {
&SaSchemaMigration{},
&SaUser{},
&SaProject{},
&SaDatasetSource{},
&SaDatasetItem{},
&SaDatasetCron{},
&SaInboxItem{},
&SaInboxSuggestion{},
&SaTask{},