diff --git a/docker-compose.yml b/docker-compose.yml --- a/docker-compose.yml +++ b/docker-compose.yml @@ -181,7 +181,6 @@ SPINDLE_MICROVM_PIPELINES_IMAGE_DIR: /var/lib/spindle/images SPINDLE_MICROVM_PIPELINES_OVERLAY_DIR: /var/lib/spindle/overlays SPINDLE_MICROVM_PIPELINES_AGENT_PORT: "11240" - SPINDLE_S3_LOG_BUCKET: "" SPINDLE_MICROVM_PIPELINES_ENABLE_CGROUPS: "false" # route guest nix substitution + uploads through the local ncps cache. # ncps re-signs on serve with cache.local's key, so the guest trusts the diff --git a/nix/vm.nix b/nix/vm.nix --- a/nix/vm.nix +++ b/nix/vm.nix @@ -162,8 +162,9 @@ }; }; + artifactStores.s3.bucket = envVarOr "SPINDLE_ARTIFACT_STORES_S3_BUCKET" "tangled-logs"; + pipelines = { - logBucket = envVarOr "SPINDLE_S3_LOG_BUCKET" ""; microvm.enableKVM = nestedVirt; nixCache = { readUrls = ["http://127.0.0.1:8501"]; diff --git a/spindle/server.go b/spindle/server.go --- a/spindle/server.go +++ b/spindle/server.go @@ -32,6 +32,7 @@ "tangled.org/core/rbac" "tangled.org/core/repoident" "tangled.org/core/repoverify" + "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/engine" @@ -71,6 +72,8 @@ motdMu sync.RWMutex rootCtx context.Context jobWake chan struct{} + stores *artifactstore.Stores + reader artifactstore.Reader } // New creates a new Spindle server with the provided configuration and engines. @@ -163,6 +166,19 @@ motd: defaultMotd, rootCtx: ctx, jobWake: make(chan struct{}, 1), + } + diskFallback := cfg.Server.LogDir + if cfg.ArtifactStores.Disk.Dir == "" { + logger.Warn("using SPINDLE_SERVER_LOG_DIR as the implicit disk artifact store; configure SPINDLE_ARTIFACT_STORES_DISK_DIR explicitly") + } + stores, err := artifactstore.NewStores(cfg.ArtifactStores, diskFallback, cfg.LegacyS3.LogBucket) + if err != nil { + return nil, fmt.Errorf("failed to setup artifact stores: %w", err) + } + spindle.stores = stores + spindle.reader = stores + if cfg.LegacyS3.LogBucket != "" { + logger.Warn("SPINDLE_S3_LOG_BUCKET is deprecated; use SPINDLE_ARTIFACT_STORES_S3_BUCKET") } err = e.AddSpindle(rbacDomain) @@ -387,16 +403,17 @@ l := log.SubLogger(s.l, "xrpc") x := xrpc.Xrpc{ - Logger: l, - Db: s.db, - Enforcer: s.e, - Engines: s.engs, - Config: s.cfg, - Resolver: s.res, - Vault: s.vault, - Notifier: s.Notifier(), - ServiceAuth: serviceAuth, - Trigger: s, + Logger: l, + Db: s.db, + Enforcer: s.e, + Engines: s.engs, + Config: s.cfg, + ArtifactReader: s.reader, + Resolver: s.res, + Vault: s.vault, + Notifier: s.Notifier(), + ServiceAuth: serviceAuth, + Trigger: s, } return x.Router() @@ -820,7 +837,7 @@ workflows[eng] = append(workflows[eng], *ewf) } - engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.db, s.n, s.rootCtx, &models.Pipeline{ + engine.StartWorkflows(log.SubLogger(s.l, "engine"), s.vault, s.cfg, s.stores, s.db, s.n, s.rootCtx, &models.Pipeline{ RepoDid: syntax.DID(job.RepoDid), Workflows: workflows, TrustedSource: trustedSource, diff --git a/spindle/stream.go b/spindle/stream.go --- a/spindle/stream.go +++ b/spindle/stream.go @@ -4,17 +4,16 @@ "context" "errors" "fmt" - "io" "net/http" "time" "tangled.org/core/eventstream" "tangled.org/core/log" + "tangled.org/core/spindle/logview" "tangled.org/core/spindle/models" "github.com/go-chi/chi/v5" "github.com/gorilla/websocket" - "github.com/hpcloud/tail" ) var upgrader = websocket.Upgrader{ @@ -89,40 +88,27 @@ } isFinished := models.StatusKind(status.Status).IsFinish() - filePath := models.LogFilePath(s.cfg.Server.LogDir, wid) - - config := tail.Config{ - Follow: !isFinished, - ReOpen: !isFinished, - MustExist: false, - Location: &tail.SeekInfo{ - Offset: 0, - Whence: io.SeekStart, - }, - // Logger: tail.DiscardingLogger, - } - - t, err := tail.TailFile(filePath, config) + lines, stop, err := logview.Follow(ctx, s.db, s.reader, s.cfg.Server.LogDir, wid, isFinished) if err != nil { - return fmt.Errorf("failed to tail log file: %w", err) + return fmt.Errorf("failed to follow workflow log: %w", err) } - defer t.Stop() - + defer stop() for { select { case <-ctx.Done(): return ctx.Err() - case line := <-t.Lines: - if line == nil && isFinished { - return fmt.Errorf("tail completed") + case line, ok := <-lines: + if !ok && isFinished { + return fmt.Errorf("log completed") } - + if !ok { + return fmt.Errorf("log channel closed unexpectedly") + } if line == nil { - return fmt.Errorf("tail channel closed unexpectedly") + continue } - if line.Err != nil { - return fmt.Errorf("error tailing log file: %w", line.Err) + return fmt.Errorf("error following workflow log: %w", line.Err) } if err := conn.WriteMessage(websocket.TextMessage, []byte(line.Text)); err != nil { diff --git a/nix/modules/spindle.nix b/nix/modules/spindle.nix --- a/nix/modules/spindle.nix +++ b/nix/modules/spindle.nix @@ -136,12 +136,27 @@ }; }; - pipelines = { - logBucket = mkOption { + artifactStores = { + disk.dir = mkOption { + type = types.path; + default = "/var/log/spindle"; + description = "Root directory for disk artifacts"; + }; + + s3.bucket = mkOption { type = types.str; default = "tangled-logs"; - description = "S3 bucket for workflow logs"; + description = "S3 bucket for artifacts"; }; + + s3.region = mkOption { + type = types.str; + default = "us-east-1"; + description = "AWS region for the artifact bucket"; + }; + }; + + pipelines = { workflowTimeout = mkOption { type = types.str; default = "5m"; @@ -388,7 +403,9 @@ "SPINDLE_NIX_CACHE_READ_URLS=${concatStringsSep "," cfg.pipelines.nixCache.readUrls}" "SPINDLE_NIX_CACHE_TRUSTED_PUBLIC_KEYS=${concatStringsSep "," cfg.pipelines.nixCache.trustedPublicKeys}" "SPINDLE_NIX_CACHE_UPLOAD_URL=${cfg.pipelines.nixCache.uploadUrl}" - "SPINDLE_S3_LOG_BUCKET=${cfg.pipelines.logBucket}" + "SPINDLE_ARTIFACT_STORES_DISK_DIR=${cfg.artifactStores.disk.dir}" + "SPINDLE_ARTIFACT_STORES_S3_BUCKET=${cfg.artifactStores.s3.bucket}" + "SPINDLE_ARTIFACT_STORES_S3_REGION=${cfg.artifactStores.s3.region}" ]; ExecStart = "${cfg.package}/bin/spindle"; Restart = "always"; diff --git a/spindle/artifactstore/artifactstore.go b/spindle/artifactstore/artifactstore.go new file mode 100644 --- /dev/null +++ b/spindle/artifactstore/artifactstore.go @@ -0,0 +1,240 @@ +package artifactstore + +import ( + "context" + "errors" + "fmt" + "io" + "os" + "path/filepath" + "strings" + + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/s3" + spindleconfig "tangled.org/core/spindle/config" +) + +type Writer interface { + Put(ctx context.Context, ref string, r io.Reader) error +} + +type Reader interface { + Open(ctx context.Context, ref string) (io.ReadCloser, error) +} + +type Store interface { + Writer + Reader +} + +type DiskStore struct { + root string +} + +func NewDiskStore(root string) (*DiskStore, error) { + if root == "" { + return nil, fmt.Errorf("artifact disk directory is required") + } + return &DiskStore{root: filepath.Clean(root)}, nil +} + +func (s *DiskStore) Put(_ context.Context, ref string, r io.Reader) error { + path, err := s.resolve(ref) + if err != nil { + return err + } + if err := os.MkdirAll(filepath.Dir(path), 0755); err != nil { + return fmt.Errorf("mkdir for artifact %q: %w", path, err) + } + tmpFile, err := os.CreateTemp(filepath.Dir(path), ".tmp-artifact-*") + if err != nil { + return fmt.Errorf("create temp artifact: %w", err) + } + tmpPath := tmpFile.Name() + defer func() { + _ = tmpFile.Close() + _ = os.Remove(tmpPath) + }() + if _, err := io.Copy(tmpFile, r); err != nil { + return fmt.Errorf("write artifact content: %w", err) + } + if err := tmpFile.Sync(); err != nil { + return fmt.Errorf("sync artifact file: %w", err) + } + if err := tmpFile.Close(); err != nil { + return fmt.Errorf("close artifact file: %w", err) + } + if err := os.Rename(tmpPath, path); err != nil { + return fmt.Errorf("rename artifact file to target: %w", err) + } + return nil +} + +func (s *DiskStore) Open(_ context.Context, ref string) (io.ReadCloser, error) { + path, err := s.resolve(ref) + if err != nil { + return nil, err + } + f, err := os.Open(path) + if err != nil { + return nil, fmt.Errorf("open disk artifact %q: %w", path, err) + } + return f, nil +} + +func (s *DiskStore) resolve(ref string) (string, error) { + if ref == "" || filepath.IsAbs(ref) { + return "", fmt.Errorf("invalid disk artifact ref %q", ref) + } + path := filepath.Join(s.root, filepath.Clean(ref)) + rel, err := filepath.Rel(s.root, path) + if err != nil || rel == ".." || strings.HasPrefix(rel, ".."+string(filepath.Separator)) { + return "", fmt.Errorf("artifact ref %q escapes disk root %q", ref, s.root) + } + return path, nil +} + +type s3API interface { + PutObject(ctx context.Context, params *s3.PutObjectInput, optFns ...func(*s3.Options)) (*s3.PutObjectOutput, error) + GetObject(ctx context.Context, params *s3.GetObjectInput, optFns ...func(*s3.Options)) (*s3.GetObjectOutput, error) +} + +type S3Store struct { + client s3API + bucket string +} + +func NewS3Store(client s3API, bucket string) (*S3Store, error) { + if client == nil { + return nil, fmt.Errorf("s3 client is required") + } + if bucket == "" { + return nil, fmt.Errorf("artifact S3 bucket is required") + } + return &S3Store{client: client, bucket: bucket}, nil +} + +func (s *S3Store) Put(ctx context.Context, ref string, r io.Reader) error { + if err := validateObjectRef(ref); err != nil { + return err + } + _, err := s.client.PutObject(ctx, &s3.PutObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(ref), + Body: r, + }) + if err != nil { + return fmt.Errorf("s3 put object: %w", err) + } + return nil +} + +func (s *S3Store) Open(ctx context.Context, ref string) (io.ReadCloser, error) { + if err := validateObjectRef(ref); err != nil { + return nil, err + } + res, err := s.client.GetObject(ctx, &s3.GetObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(ref), + }) + if err != nil { + return nil, fmt.Errorf("s3 get object: %w", err) + } + return res.Body, nil +} + +func validateObjectRef(ref string) error { + if ref == "" || strings.HasPrefix(ref, "/") || strings.Contains(ref, "://") { + return fmt.Errorf("invalid artifact ref %q", ref) + } + return nil +} + +type Stores struct { + order []string + stores map[string]Store +} + +func NewStores(cfg spindleconfig.ArtifactStores, diskFallback, legacyS3Bucket string) (*Stores, error) { + stores := &Stores{stores: make(map[string]Store)} + diskDir := cfg.Disk.Dir + if diskDir == "" { + diskDir = diskFallback + } + if diskDir != "" { + disk, err := NewDiskStore(diskDir) + if err != nil { + return nil, err + } + stores.order = append(stores.order, "disk") + stores.stores["disk"] = disk + } + + bucket := cfg.S3.Bucket + if bucket == "" { + bucket = legacyS3Bucket + } + if bucket != "" { + awsCfg, err := config.LoadDefaultConfig(context.Background(), config.WithRegion(cfg.S3.Region)) + if err != nil { + return nil, fmt.Errorf("load aws config: %w", err) + } + s3Store, err := NewS3Store(s3.NewFromConfig(awsCfg), bucket) + if err != nil { + return nil, err + } + stores.order = append(stores.order, "s3") + stores.stores["s3"] = s3Store + } + return stores, nil +} + +func (s *Stores) Names() []string { + return append([]string(nil), s.order...) +} + +func (s *Stores) Store(name string) (Store, bool) { + store, ok := s.stores[name] + return store, ok +} + +func (s *Stores) Open(ctx context.Context, ref string) (io.ReadCloser, error) { + var errs []error + for _, name := range s.order { + rc, err := s.stores[name].Open(ctx, ref) + if err == nil { + return rc, nil + } + errs = append(errs, fmt.Errorf("%s: %w", name, err)) + } + return nil, fmt.Errorf("open artifact %q: %w", ref, errors.Join(errs...)) +} + +func (s *Stores) PutFile(ctx context.Context, ref, sourcePath string) []error { + var errs []error + for _, name := range s.order { + store := s.stores[name] + if disk, ok := store.(*DiskStore); ok { + target, err := disk.resolve(ref) + if err == nil { + source, sourceErr := filepath.Abs(sourcePath) + targetAbs, targetErr := filepath.Abs(target) + if sourceErr == nil && targetErr == nil && source == targetAbs { + continue + } + } + } + file, err := os.Open(sourcePath) + if err != nil { + errs = append(errs, fmt.Errorf("%s: open source: %w", name, err)) + continue + } + err = store.Put(ctx, ref, file) + _ = file.Close() + if err != nil { + errs = append(errs, fmt.Errorf("%s: %w", name, err)) + } + } + return errs +} diff --git a/spindle/artifactstore/artifactstore_test.go b/spindle/artifactstore/artifactstore_test.go new file mode 100644 --- /dev/null +++ b/spindle/artifactstore/artifactstore_test.go @@ -0,0 +1,130 @@ +package artifactstore + +import ( + "bytes" + "context" + "io" + "os" + "strings" + "sync" + "testing" + + "github.com/aws/aws-sdk-go-v2/service/s3" +) + +type mockS3Client struct { + mu sync.Mutex + store map[string][]byte +} + +func newMockS3Client() *mockS3Client { + return &mockS3Client{ + store: make(map[string][]byte), + } +} + +func (m *mockS3Client) PutObject(ctx context.Context, params *s3.PutObjectInput, optFns ...func(*s3.Options)) (*s3.PutObjectOutput, error) { + m.mu.Lock() + defer m.mu.Unlock() + b, err := io.ReadAll(params.Body) + if err != nil { + return nil, err + } + key := *params.Bucket + "/" + *params.Key + m.store[key] = b + return &s3.PutObjectOutput{}, nil +} + +func (m *mockS3Client) GetObject(ctx context.Context, params *s3.GetObjectInput, optFns ...func(*s3.Options)) (*s3.GetObjectOutput, error) { + m.mu.Lock() + defer m.mu.Unlock() + key := *params.Bucket + "/" + *params.Key + data, ok := m.store[key] + if !ok { + return nil, os.ErrNotExist + } + return &s3.GetObjectOutput{ + Body: io.NopCloser(bytes.NewReader(data)), + }, nil +} + +func TestDiskStore(t *testing.T) { + tempDir := t.TempDir() + store, err := NewDiskStore(tempDir) + if err != nil { + t.Fatal(err) + } + + ctx := context.Background() + ref := "logs/test.log" + content := "hello world log content" + + if err := store.Put(ctx, ref, strings.NewReader(content)); err != nil { + t.Fatalf("Put failed: %v", err) + } + + rc, err := store.Open(ctx, ref) + if err != nil { + t.Fatalf("Open failed: %v", err) + } + defer rc.Close() + + got, err := io.ReadAll(rc) + if err != nil { + t.Fatalf("ReadAll failed: %v", err) + } + if string(got) != content { + t.Fatalf("got content %q, want %q", string(got), content) + } +} + +func TestDiskStoreTraversalProtection(t *testing.T) { + tempDir := t.TempDir() + store, err := NewDiskStore(tempDir) + if err != nil { + t.Fatal(err) + } + + ctx := context.Background() + badRef := "../outside" + + err = store.Put(ctx, badRef, strings.NewReader("bad")) + if err == nil { + t.Fatal("expected error putting file outside diskDir, got nil") + } + + _, err = store.Open(ctx, badRef) + if err == nil { + t.Fatal("expected error opening file outside diskDir, got nil") + } +} + +func TestS3Store(t *testing.T) { + mock := newMockS3Client() + store, err := NewS3Store(mock, "mybucket") + if err != nil { + t.Fatal(err) + } + + ctx := context.Background() + ref := "logs/run1.log" + content := "s3 log payload" + + if err := store.Put(ctx, ref, strings.NewReader(content)); err != nil { + t.Fatalf("Put to S3 failed: %v", err) + } + + rc, err := store.Open(ctx, ref) + if err != nil { + t.Fatalf("Open from S3 failed: %v", err) + } + defer rc.Close() + + got, err := io.ReadAll(rc) + if err != nil { + t.Fatalf("ReadAll failed: %v", err) + } + if string(got) != content { + t.Fatalf("got content %q, want %q", string(got), content) + } +} diff --git a/spindle/config/config.go b/spindle/config/config.go --- a/spindle/config/config.go +++ b/spindle/config/config.go @@ -57,7 +57,21 @@ MaxConcurrentWorkflows int `env:"MAX_CONCURRENT_WORKFLOWS, default=8"` // max number of workflow containers running at once (memory cap) } -type S3 struct { +type ArtifactStoreDisk struct { + Dir string `env:"DIR"` +} + +type ArtifactStoreS3 struct { + Bucket string `env:"BUCKET"` + Region string `env:"REGION, default=us-east-1"` +} + +type ArtifactStores struct { + Disk ArtifactStoreDisk `env:",prefix=DISK_"` + S3 ArtifactStoreS3 `env:",prefix=S3_"` +} + +type LegacyS3 struct { LogBucket string `env:"LOG_BUCKET"` } @@ -98,7 +112,8 @@ NixeryPipelines NixeryPipelines `env:",prefix=SPINDLE_NIXERY_PIPELINES_"` MicroVMPipelines MicroVMPipelines `env:",prefix=SPINDLE_MICROVM_PIPELINES_"` NixCache NixCache `env:",prefix=SPINDLE_NIX_CACHE_"` - S3 S3 `env:",prefix=SPINDLE_S3_"` + ArtifactStores ArtifactStores `env:",prefix=SPINDLE_ARTIFACT_STORES_"` + LegacyS3 LegacyS3 `env:",prefix=SPINDLE_S3_"` } func Load(ctx context.Context) (*Config, error) { diff --git a/spindle/db/artifacts.go b/spindle/db/artifacts.go new file mode 100644 --- /dev/null +++ b/spindle/db/artifacts.go @@ -0,0 +1,32 @@ +package db + +type FinishedLog struct { + LeaseID string + Workflow string + Ref string + Hash string +} + +func (d *DB) GetFinishedLog(workflow string) (*FinishedLog, error) { + var fl FinishedLog + err := d.QueryRow( + `select lease_id, workflow, ref, hash + from mill_artifacts + where workflow = ? + order by id desc limit 1`, + workflow, + ).Scan(&fl.LeaseID, &fl.Workflow, &fl.Ref, &fl.Hash) + if err != nil { + return nil, err + } + return &fl, nil +} + +func (d *DB) SaveArtifactRef(leaseID, workflow, ref, hash string) error { + _, err := d.Exec( + `insert into mill_artifacts (lease_id, workflow, ref, hash) + values (?, ?, ?, ?)`, + leaseID, workflow, ref, hash, + ) + return err +} diff --git a/spindle/db/db.go b/spindle/db/db.go --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -132,6 +132,14 @@ foreign key (pipeline_id) references pipelines(id) on delete cascade ); + create table if not exists mill_artifacts ( + id integer primary key autoincrement, + lease_id text not null, + workflow text not null, + ref text not null, + hash text not null + ); + create table if not exists migrations ( id integer primary key autoincrement, name text unique diff --git a/spindle/engine/engine.go b/spindle/engine/engine.go --- a/spindle/engine/engine.go +++ b/spindle/engine/engine.go @@ -2,13 +2,18 @@ import ( "context" + "crypto/sha256" + "encoding/hex" "errors" "fmt" + "io" "log/slog" - "path/filepath" + "os" "sync" + "time" "tangled.org/core/notifier" + "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" @@ -66,7 +71,7 @@ FinalizeWorkflow(ctx context.Context, wid models.WorkflowId, wf *models.Workflow, wfLogger models.WorkflowLogger) error } -func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, db *db.DB, n *notifier.Notifier, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { +func StartWorkflows(l *slog.Logger, vault secrets.Manager, cfg *config.Config, stores *artifactstore.Stores, db *db.DB, n *notifier.Notifier, ctx context.Context, pipeline *models.Pipeline, pipelineId models.PipelineId) { l.Info("starting all workflows in parallel", "pipeline", pipelineId) var allSecrets []secrets.UnlockedSecret @@ -82,11 +87,6 @@ secretValues := make([]string, len(allSecrets)) for i, s := range allSecrets { secretValues[i] = s.Value - } - - s3, err := NewS3(cfg.S3.LogBucket) - if err != nil { - l.Error("error creating s3 client", "err", err) } // wid.String() is lossy so two different names can map to the same key @@ -128,21 +128,13 @@ l.Info("skipping finished workflow", "wid", wid, "status", st.Status) return } - defer func() { - if s3 != nil { - logFile := filepath.Join(cfg.Server.LogDir, fmt.Sprintf("%s.log", wid.String())) - if err := s3.WriteFile(ctx, logFile); err != nil { - l.Error("error uploading logs", "err", err) - } - } - }() - wfLogger, err := models.NewFileWorkflowLogger(cfg.Server.LogDir, wid, secretValues) if err != nil { l.Warn("failed to setup step logger; logs will not be persisted", "error", err) wfLogger = models.NullLogger{} } else { l.Info("setup step logger; logs will be persisted", "logDir", cfg.Server.LogDir, "wid", wid) + defer archiveWorkflowLog(l, stores, db, cfg.Server.LogDir, wid) defer wfLogger.Close() } @@ -234,4 +226,38 @@ wg.Wait() l.Info("all workflows completed") +} + +func archiveWorkflowLog(l *slog.Logger, stores *artifactstore.Stores, database *db.DB, logDir string, wid models.WorkflowId) { + if stores == nil { + return + } + logPath := models.LogFilePath(logDir, wid) + file, err := os.Open(logPath) + if err != nil { + l.Error("open workflow log for archival", "wid", wid, "err", err) + return + } + hash := sha256.New() + if _, err := io.Copy(hash, file); err != nil { + _ = file.Close() + l.Error("hash workflow log", "wid", wid, "err", err) + return + } + _ = file.Close() + + ref := wid.String() + ".log" + uploadCtx, cancel := context.WithTimeout(context.Background(), 2*time.Minute) + defer cancel() + errs := stores.PutFile(uploadCtx, ref, logPath) + for _, err := range errs { + l.Error("archive workflow log", "wid", wid, "err", err) + } + if len(errs) == len(stores.Names()) { + return + } + digest := "sha256:" + hex.EncodeToString(hash.Sum(nil)) + if err := database.SaveArtifactRef(wid.String(), wid.Name, ref, digest); err != nil { + l.Error("save workflow log artifact", "wid", wid, "err", err) + } } diff --git a/spindle/engine/engine_test.go b/spindle/engine/engine_test.go --- a/spindle/engine/engine_test.go +++ b/spindle/engine/engine_test.go @@ -114,7 +114,7 @@ } cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}} - StartWorkflows(logger, nil, cfg, testDB, nil, context.Background(), pipeline, pipelineId) + StartWorkflows(logger, nil, cfg, nil, testDB, nil, context.Background(), pipeline, pipelineId) eng.mu.Lock() setupCalls := append([]models.WorkflowId(nil), eng.setupCalls...) @@ -195,7 +195,7 @@ cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}} doneChan := make(chan struct{}) go func() { - StartWorkflows(logger, nil, cfg, testDB, nil, context.Background(), pipeline, pipelineId) + StartWorkflows(logger, nil, cfg, nil, testDB, nil, context.Background(), pipeline, pipelineId) close(doneChan) }() @@ -250,7 +250,7 @@ } cfg := &config.Config{Server: config.Server{LogDir: t.TempDir()}} - StartWorkflows(logger, nil, cfg, testDB, nil, context.Background(), pipeline, pipelineId) + StartWorkflows(logger, nil, cfg, nil, testDB, nil, context.Background(), pipeline, pipelineId) st, err := testDB.GetStatus(wid) if err != nil { diff --git a/spindle/engine/s3.go b/spindle/engine/s3.go deleted file mode 100644 --- a/spindle/engine/s3.go +++ /dev/null @@ -1,75 +0,0 @@ -package engine - -import ( - "context" - "fmt" - "io" - "os" - "path/filepath" - - "github.com/aws/aws-sdk-go-v2/aws" - "github.com/aws/aws-sdk-go-v2/config" - "github.com/aws/aws-sdk-go-v2/service/s3" -) - -type S3 struct { - bucket string - client *s3.Client -} - -const BaseS3Path = "spindle/workflows" - -func NewS3(bucket string) (*S3, error) { - if bucket == "" { - return nil, fmt.Errorf("s3 bucket not provided") - } - - ctx := context.Background() - sdkConfig, err := config.LoadDefaultConfig(ctx) - - if err != nil { - return nil, fmt.Errorf("error loading s3 config: %w", err) - } - s3Client := s3.NewFromConfig(sdkConfig) - - return &S3{ - bucket: bucket, - client: s3Client, - }, nil -} - -func (s *S3) WriteFile(ctx context.Context, path string) error { - s3Key := fmt.Sprintf("%s/%s", BaseS3Path, filepath.Base(path)) - - file, err := os.Open(path) - if err != nil { - return fmt.Errorf("error opening file %s: %w", path, err) - } - defer file.Close() - - _, err = s.client.PutObject(ctx, &s3.PutObjectInput{ - Bucket: &s.bucket, - Key: &s3Key, - Body: file, - }) - - if err != nil { - return fmt.Errorf("error writing to s3: %w", err) - } - - return nil -} - -func (s *S3) ReadFile(ctx context.Context, name string) ([]byte, error) { - res, err := s.client.GetObject(ctx, &s3.GetObjectInput{ - Bucket: &s.bucket, - Key: aws.String(name), - }) - - if err != nil { - return nil, fmt.Errorf("error reading file %s: %w", name, err) - } - defer res.Body.Close() - - return io.ReadAll(res.Body) -} diff --git a/spindle/logview/logview.go b/spindle/logview/logview.go new file mode 100644 --- /dev/null +++ b/spindle/logview/logview.go @@ -0,0 +1,52 @@ +package logview + +import ( + "bufio" + "context" + "io" + "strings" + + "github.com/hpcloud/tail" + "tangled.org/core/spindle/artifactstore" + "tangled.org/core/spindle/db" + "tangled.org/core/spindle/models" +) + +// streams a workflow's log lines: tails the log file while running, reads +// the uploaded artifact once finished, same for every role. stop ends a +// live follow early, the channel closes when the source drains or ctx ends +func Follow(ctx context.Context, d *db.DB, reader artifactstore.Reader, logDir string, wid models.WorkflowId, finished bool) (<-chan *tail.Line, func(), error) { + if finished && reader != nil && d != nil { + if fl, err := d.GetFinishedLog(wid.Name); err == nil && fl.Ref != "" { + rc, err := reader.Open(ctx, fl.Ref) + if err == nil { + ch := make(chan *tail.Line, 64) + followCtx, cancel := context.WithCancel(ctx) + go func() { + defer close(ch) + defer rc.Close() + scanner := bufio.NewScanner(rc) + for scanner.Scan() { + select { + case <-followCtx.Done(): + return + case ch <- &tail.Line{Text: strings.TrimSuffix(scanner.Text(), "\r")}: + } + } + }() + return ch, cancel, nil + } + } + } + + t, err := tail.TailFile(models.LogFilePath(logDir, wid), tail.Config{ + Follow: !finished, + ReOpen: !finished, + MustExist: false, + Location: &tail.SeekInfo{Offset: 0, Whence: io.SeekStart}, + }) + if err != nil { + return nil, nil, err + } + return t.Lines, func() { _ = t.Stop() }, nil +} diff --git a/spindle/xrpc/ci_pipeline_subscribe_logs.go b/spindle/xrpc/ci_pipeline_subscribe_logs.go --- a/spindle/xrpc/ci_pipeline_subscribe_logs.go +++ b/spindle/xrpc/ci_pipeline_subscribe_logs.go @@ -4,7 +4,6 @@ "context" "encoding/json" "fmt" - "io" "net/http" "sync" "time" @@ -12,8 +11,8 @@ "github.com/bluesky-social/indigo/atproto/atclient" "github.com/bluesky-social/indigo/atproto/syntax" "github.com/gorilla/websocket" - "github.com/hpcloud/tail" "tangled.org/core/api/tangled" + "tangled.org/core/spindle/logview" "tangled.org/core/spindle/models" ) @@ -166,26 +165,14 @@ isFinished = models.StatusKind(status.Status).IsFinish() } - filePath := models.LogFilePath(x.Config.Server.LogDir, wid) - - tailConfig := tail.Config{ - Follow: !isFinished, - ReOpen: !isFinished, - MustExist: false, - Location: &tail.SeekInfo{ - Offset: 0, - Whence: io.SeekStart, - }, - } - - t, err := tail.TailFile(filePath, tailConfig) + lines, stop, err := logview.Follow(ctx, x.Db, x.ArtifactReader, x.Config.Server.LogDir, wid, isFinished) if err != nil { - l.Error("failed to tail log file", "workflow", wfName, "err", err) + l.Error("failed to follow workflow log", "workflow", wfName, "err", err) return } - defer t.Stop() + defer stop() - // if we are following, poll status in database to stop tailing when finished + // if we are following, poll status in database to stop when finished if !isFinished { go func() { ticker := time.NewTicker(2 * time.Second) @@ -197,7 +184,7 @@ case <-ticker.C: status, err := x.Db.GetStatus(wid) if err == nil && models.StatusKind(status.Status).IsFinish() { - t.Stop() + stop() return } } @@ -209,7 +196,7 @@ select { case <-ctx.Done(): return - case line, ok := <-t.Lines: + case line, ok := <-lines: if !ok || line == nil { return } diff --git a/spindle/xrpc/xrpc.go b/spindle/xrpc/xrpc.go --- a/spindle/xrpc/xrpc.go +++ b/spindle/xrpc/xrpc.go @@ -15,6 +15,7 @@ "tangled.org/core/idresolver" "tangled.org/core/notifier" "tangled.org/core/rbac" + "tangled.org/core/spindle/artifactstore" "tangled.org/core/spindle/config" "tangled.org/core/spindle/db" "tangled.org/core/spindle/models" @@ -41,16 +42,17 @@ } type Xrpc struct { - Logger *slog.Logger - Db *db.DB - Enforcer *rbac.Enforcer - Engines map[string]models.Engine - Config *config.Config - Resolver *idresolver.Resolver - Vault secrets.Manager - Notifier *notifier.Notifier - ServiceAuth *serviceauth.ServiceAuth - Trigger PipelineTrigger + Logger *slog.Logger + Db *db.DB + Enforcer *rbac.Enforcer + Engines map[string]models.Engine + Config *config.Config + ArtifactReader artifactstore.Reader + Resolver *idresolver.Resolver + Vault secrets.Manager + Notifier *notifier.Notifier + ServiceAuth *serviceauth.ServiceAuth + Trigger PipelineTrigger } func (x *Xrpc) Router() http.Handler {