diff --git a/nodes/node.go b/nodes/node.go index c71bca8..75f1663 100644 --- a/nodes/node.go +++ b/nodes/node.go @@ -59,3 +59,7 @@ func NewContext(executionID, pipeID string, db *store.DB) *Context { func (c *Context) Log(nodeID, level, message string) { c.DB.LogExecution(c.ExecutionID, nodeID, level, message) } + +func (c *Context) SaveOutput(format, content, contentType string) error { + return c.DB.SavePipeOutput(c.PipeID, format, content, contentType) +} diff --git a/nodes/outputs/json.go b/nodes/outputs/json.go index 7d42723..f6096c7 100644 --- a/nodes/outputs/json.go +++ b/nodes/outputs/json.go @@ -33,9 +33,13 @@ func (n *JSONOutputNode) Execute(ctx context.Context, config map[string]interfac return nil, err } + // Save output to database for public access + if err := execCtx.SaveOutput("json", string(jsonData), "application/json"); err != nil { + execCtx.Log("json-output", "error", "Failed to save output: "+err.Error()) + } + execCtx.Log("json-output", "info", string(jsonData)) - // Return the data (for potential chaining) return data, nil } diff --git a/nodes/outputs/rss.go b/nodes/outputs/rss.go index 9bf9e71..2cc1bce 100644 --- a/nodes/outputs/rss.go +++ b/nodes/outputs/rss.go @@ -87,6 +87,12 @@ func (n *RSSOutputNode) Execute(ctx context.Context, config map[string]interface } rssOutput := fmt.Sprintf("\n%s", string(xmlData)) + + // Save output to database for public access + if err := execCtx.SaveOutput("rss", rssOutput, "application/rss+xml"); err != nil { + execCtx.Log("rss-output", "error", "Failed to save output: "+err.Error()) + } + execCtx.Log("rss-output", "info", rssOutput) return data, nil diff --git a/store/db.go b/store/db.go index 62ba1b7..a423457 100644 --- a/store/db.go +++ b/store/db.go @@ -78,6 +78,19 @@ func (db *DB) initSchema() error { CREATE INDEX IF NOT EXISTS idx_pipes_user_id ON pipes(user_id); + -- Pipe outputs (cached output for public feeds) + CREATE TABLE IF NOT EXISTS pipe_outputs ( + id TEXT PRIMARY KEY, + pipe_id TEXT NOT NULL REFERENCES pipes(id) ON DELETE CASCADE, + format TEXT NOT NULL, + content TEXT NOT NULL, + content_type TEXT NOT NULL, + created_at INTEGER NOT NULL + ); + + CREATE INDEX IF NOT EXISTS idx_outputs_pipe_id ON pipe_outputs(pipe_id); + CREATE UNIQUE INDEX IF NOT EXISTS idx_outputs_pipe_format ON pipe_outputs(pipe_id, format); + -- Scheduled jobs CREATE TABLE IF NOT EXISTS scheduled_jobs ( id TEXT PRIMARY KEY, diff --git a/store/pipes.go b/store/pipes.go index 76f498b..2969139 100644 --- a/store/pipes.go +++ b/store/pipes.go @@ -19,6 +19,15 @@ type Pipe struct { UpdatedAt int64 `json:"updated_at"` } +type PipeOutput struct { + ID string `json:"id"` + PipeID string `json:"pipe_id"` + Format string `json:"format"` + Content string `json:"content"` + ContentType string `json:"content_type"` + CreatedAt int64 `json:"created_at"` +} + type ScheduledJob struct { ID string PipeID string @@ -130,6 +139,46 @@ func (db *DB) DeletePipe(id string) error { return nil } +func (db *DB) SavePipeOutput(pipeID, format, content, contentType string) error { + now := time.Now().Unix() + id := uuid.New().String() + + _, err := db.Exec(` + INSERT INTO pipe_outputs (id, pipe_id, format, content, content_type, created_at) + VALUES (?, ?, ?, ?, ?, ?) + ON CONFLICT(pipe_id, format) DO UPDATE SET + content = excluded.content, + content_type = excluded.content_type, + created_at = excluded.created_at + `, id, pipeID, format, content, contentType, now) + + if err != nil { + return fmt.Errorf("save pipe output: %w", err) + } + + return nil +} + +func (db *DB) GetPipeOutput(pipeID, format string) (*PipeOutput, error) { + output := &PipeOutput{} + + err := db.QueryRow(` + SELECT id, pipe_id, format, content, content_type, created_at + FROM pipe_outputs + WHERE pipe_id = ? AND format = ? + `, pipeID, format).Scan(&output.ID, &output.PipeID, &output.Format, &output.Content, &output.ContentType, &output.CreatedAt) + + if err == sql.ErrNoRows { + return nil, nil + } + + if err != nil { + return nil, fmt.Errorf("get pipe output: %w", err) + } + + return output, nil +} + func (db *DB) CreateScheduledJob(pipeID, cronExpression string, nextRunAt int64) (*ScheduledJob, error) { now := time.Now().Unix() job := &ScheduledJob{ diff --git a/web/server.go b/web/server.go index d80305c..f6a034a 100644 --- a/web/server.go +++ b/web/server.go @@ -7,6 +7,7 @@ import ( "html/template" "net/http" "strconv" + "strings" "github.com/charmbracelet/log" "github.com/kierank/pipes/auth" @@ -68,6 +69,9 @@ func (s *Server) Start() error { mux.HandleFunc("/api/node-types", s.handleAPINodeTypes) mux.HandleFunc("/api/executions/", s.sessionManager.RequireAuth(s.handleAPIExecution)) + // Public feed routes + mux.HandleFunc("/feeds/", s.handlePublicFeed) + s.server = &http.Server{ Addr: fmt.Sprintf("%s:%d", s.cfg.Host, s.cfg.Port), Handler: mux, @@ -322,6 +326,7 @@ func (s *Server) handleAPIPipe(w http.ResponseWriter, r *http.Request) { Name string `json:"name"` Description string `json:"description"` Config map[string]interface{} `json:"config"` + IsPublic *bool `json:"is_public"` } if err := json.NewDecoder(r.Body).Decode(&req); err != nil { @@ -339,6 +344,9 @@ func (s *Server) handleAPIPipe(w http.ResponseWriter, r *http.Request) { configJSON, _ := json.Marshal(req.Config) pipe.Config = string(configJSON) } + if req.IsPublic != nil { + pipe.IsPublic = *req.IsPublic + } if err := s.db.UpdatePipe(pipe); err != nil { http.Error(w, "Failed to update pipe", http.StatusInternalServerError) @@ -520,6 +528,76 @@ func (s *Server) handleExecutionLogs(w http.ResponseWriter, r *http.Request, exe json.NewEncoder(w).Encode(logs) } +func (s *Server) handlePublicFeed(w http.ResponseWriter, r *http.Request) { + // Parse path: /feeds/{id}.{format} or /feeds/{id}/{format} + path := strings.TrimPrefix(r.URL.Path, "/feeds/") + if path == "" { + http.Error(w, "Not found", http.StatusNotFound) + return + } + + var pipeID, format string + + // Check for extension format: id.json or id.rss + if strings.Contains(path, ".") { + parts := strings.SplitN(path, ".", 2) + pipeID = parts[0] + format = parts[1] + } else if strings.Contains(path, "/") { + // Check for path format: id/json or id/rss + parts := strings.SplitN(path, "/", 2) + pipeID = parts[0] + format = parts[1] + } else { + // Default to json if no format specified + pipeID = path + format = "json" + } + + // Look up pipe by ID + pipe, err := s.db.GetPipe(pipeID) + if err != nil { + s.logger.Error("failed to get pipe", "pipe_id", pipeID, "error", err) + http.Error(w, "Internal server error", http.StatusInternalServerError) + return + } + + if pipe == nil || !pipe.IsPublic { + http.Error(w, "Feed not found", http.StatusNotFound) + return + } + + // Get the cached output + output, err := s.db.GetPipeOutput(pipe.ID, format) + if err != nil { + s.logger.Error("failed to get pipe output", "pipe_id", pipe.ID, "format", format, "error", err) + http.Error(w, "Internal server error", http.StatusInternalServerError) + return + } + + // Auto-run if no output exists + if output == nil { + executor := engine.NewExecutor(s.db) + _, err := executor.Execute(r.Context(), pipe.ID, "auto") + if err != nil { + s.logger.Error("auto-execute failed", "pipe_id", pipe.ID, "error", err) + http.Error(w, "Failed to generate feed", http.StatusInternalServerError) + return + } + + // Try to get output again + output, err = s.db.GetPipeOutput(pipe.ID, format) + if err != nil || output == nil { + http.Error(w, "Feed not available in requested format", http.StatusNotFound) + return + } + } + + w.Header().Set("Content-Type", output.ContentType) + w.Header().Set("Cache-Control", "public, max-age=300") + w.Write([]byte(output.Content)) +} + // Helper functions func (s *Server) renderError(w http.ResponseWriter, title, message, details string) { diff --git a/web/templates/editor.html b/web/templates/editor.html index 2c311fe..8d74ba8 100644 --- a/web/templates/editor.html +++ b/web/templates/editor.html @@ -373,6 +373,10 @@ {{.Pipe.Name}}