package app import ( "context" "errors" "fmt" "log/slog" "net/http" "os" "strings" "time" "github.com/tim/cairnquire/apps/server/internal/auth" "github.com/tim/cairnquire/apps/server/internal/collaboration" "github.com/tim/cairnquire/apps/server/internal/config" "github.com/tim/cairnquire/apps/server/internal/database" "github.com/tim/cairnquire/apps/server/internal/docs" "github.com/tim/cairnquire/apps/server/internal/email" "github.com/tim/cairnquire/apps/server/internal/httpserver" "github.com/tim/cairnquire/apps/server/internal/markdown" "github.com/tim/cairnquire/apps/server/internal/realtime" "github.com/tim/cairnquire/apps/server/internal/store" "github.com/tim/cairnquire/apps/server/internal/sync" ) type App struct { cfg config.Config logger *slog.Logger db database.DB docs *docs.Service hub *realtime.Hub server *http.Server emailQueue *email.Queue } func New(ctx context.Context, cfg config.Config, logger *slog.Logger) (*App, error) { if err := os.MkdirAll(cfg.Content.SourceDir, 0o755); err != nil { return nil, fmt.Errorf("create content source dir: %w", err) } contentStore, err := store.New(cfg.Content.StoreDir) if err != nil { return nil, fmt.Errorf("create content store: %w", err) } db, err := database.Open(ctx, cfg.Database) if err != nil { return nil, fmt.Errorf("open database: %w", err) } if err := database.ApplyMigrations(ctx, db.SQL()); err != nil { return nil, fmt.Errorf("apply migrations: %w", err) } renderer := markdown.NewRenderer() repo := docs.NewRepository(db.SQL()) service := docs.NewService(cfg.Content.SourceDir, contentStore, renderer, repo, logger) hub := realtime.NewHub(logger) if _, err := service.SyncSourceDir(ctx); err != nil { logger.Warn("initial content sync failed", "error", err) } syncRepo := sync.NewRepository(db.SQL()) syncService := sync.NewService(syncRepo, service, contentStore, cfg.Content.SourceDir, logger) authRepo := auth.NewRepository(db.SQL()) authService, err := auth.NewService(authRepo, cfg.Auth.PublicOrigin) if err != nil { return nil, fmt.Errorf("build auth service: %w", err) } collabRepo := collaboration.NewRepository(db.SQL()) var emailSender collaboration.EmailSender var emailQueue *email.Queue emailRenderer := email.CollaborationRenderer{InstanceName: cfg.Email.InstanceName} switch cfg.Email.ResolvedProvider() { case "postmark": pm := email.NewPostmarkSender(cfg.Email.PostmarkToken, cfg.Email.From, cfg.Email.PostmarkAPIURL, logger) emailSender = pm authService.SetEmailSender(pm) case "smtp": sm := email.NewSMTPSender(cfg.Email, logger) emailSender = sm authService.SetEmailSender(sm) default: noOp := email.NewNoOpSender(logger) emailSender = noOp authService.SetEmailSender(noOp) } // The durable queue wraps whichever sender is configured so retries // survive process restarts. It is nil-safe in the service: if absent, // notifications stay in-app only. emailQueue = email.NewQueue(db.SQL(), emailSender, logger) collabService := collaboration.NewServiceWithOptions(collabRepo, hub, logger, collaboration.Options{ EmailSender: emailSender, EmailQueue: emailQueue, EmailRenderer: emailRenderer, UserLookup: authRepo, PublicOrigin: cfg.Auth.PublicOrigin, }) service.OnChange(func(change docs.DocumentChange) { hub.Broadcast(realtime.Event{Type: "document_version", Data: change}) if err := collabService.NotifyDocumentChanged(ctx, formatDocumentID(change.Path), change.Path, "system"); err != nil { logger.Warn("notify document changed", "error", err) } }) handler, err := httpserver.New(httpserver.Dependencies{ Config: cfg, Logger: logger, Documents: service, Repository: repo, ContentStore: contentStore, Hub: hub, SyncService: syncService, SyncRepo: syncRepo, Auth: authService, Collaboration: collabService, }) if err != nil { return nil, fmt.Errorf("build http handler: %w", err) } server := &http.Server{ Addr: cfg.Server.Addr, Handler: handler, ReadHeaderTimeout: 5 * time.Second, ReadTimeout: 15 * time.Second, WriteTimeout: 30 * time.Second, IdleTimeout: 60 * time.Second, } return &App{ cfg: cfg, logger: logger, db: db, docs: service, hub: hub, server: server, emailQueue: emailQueue, }, nil } func (a *App) Run(ctx context.Context) error { go a.watchDocuments(ctx) go a.syncPoll(ctx) if a.emailQueue != nil { go a.emailQueue.Run(ctx) } go func() { <-ctx.Done() shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() if err := a.server.Shutdown(shutdownCtx); err != nil && !errors.Is(err, http.ErrServerClosed) { a.logger.Error("shutdown server", "error", err) } }() a.logger.Info("starting server", "addr", a.cfg.Server.Addr) return a.server.ListenAndServe() } func (a *App) watchDocuments(ctx context.Context) { watcher, err := NewFileWatcher(a.cfg.Content.SourceDir, a.logger) if err != nil { a.logger.Warn("file watcher failed, falling back to polling only", "error", err) <-ctx.Done() return } watcher.SetOnChange(func(path string) { if _, err := a.docs.SyncSourceDir(ctx); err != nil { a.logger.Warn("document sync failed", "error", err) } else { a.logger.Debug("document synced", "path", path) } }) watcher.Run(ctx) } func (a *App) syncPoll(ctx context.Context) { ticker := time.NewTicker(30 * time.Second) defer ticker.Stop() for { select { case <-ctx.Done(): return case <-ticker.C: if _, err := a.docs.SyncSourceDir(ctx); err != nil { a.logger.Warn("document sync poll failed", "error", err) } } } } func (a *App) Close() error { return a.db.Close() } func formatDocumentID(path string) string { clean := strings.TrimPrefix(path, "/") clean = strings.TrimSuffix(clean, ".md") return "doc:" + clean }