diff --git a/cmd/flow/main.go b/cmd/flow/main.go index 7efac0c..856981a 100644 --- a/cmd/flow/main.go +++ b/cmd/flow/main.go @@ -160,10 +160,12 @@ func serve() int { } } - // Register executors now that storage is ready; the Build executor - // reads from AgentStore. + // Register executors now that storage is ready; Build reads from + // AgentStore for clone targets, Deploy looks up cloud credentials + // (agent-scoped first, global fallback). executors.RegisterAll(executors.RegistryDeps{ - Agents: mongo.Agents(), + Agents: mongo.Agents(), + Credentials: mongo.Credentials(), }) lookup = executors.BuildLookup() @@ -216,6 +218,7 @@ func serve() int { Pipelines: mongo.Pipelines(), Runs: mongo.Runs(), Agents: mongo.Agents(), + Credentials: mongo.Credentials(), Events: eventBus, Logs: logs, }), diff --git a/go.mod b/go.mod index 1f0f4ee..85bcca6 100644 --- a/go.mod +++ b/go.mod @@ -17,6 +17,22 @@ require ( require ( github.com/AdaLogics/go-fuzz-headers v0.0.0-20230811130428-ced1acdcaa24 // indirect github.com/Microsoft/go-winio v0.6.2 // indirect + github.com/aws/aws-sdk-go-v2 v1.41.7 // indirect + github.com/aws/aws-sdk-go-v2/config v1.32.17 // indirect + github.com/aws/aws-sdk-go-v2/credentials v1.19.16 // indirect + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.23 // indirect + github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.23 // indirect + github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.23 // indirect + github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.24 // indirect + github.com/aws/aws-sdk-go-v2/service/ecr v1.57.2 // indirect + github.com/aws/aws-sdk-go-v2/service/iam v1.53.10 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.9 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.23 // indirect + github.com/aws/aws-sdk-go-v2/service/signin v1.0.11 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.30.17 // indirect + github.com/aws/aws-sdk-go-v2/service/ssooidc v1.35.21 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.42.1 // indirect + github.com/aws/smithy-go v1.25.1 // indirect github.com/bahlo/generic-list-go v0.2.0 // indirect github.com/buger/jsonparser v1.1.1 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect diff --git a/go.sum b/go.sum index b3f2d1f..c598fcf 100644 --- a/go.sum +++ b/go.sum @@ -14,6 +14,38 @@ github.com/Microsoft/hcsshim v0.12.8 h1:BtDWYlFMcWhorrvSSo2M7z0csPdw6t7no/C3FsSv github.com/Microsoft/hcsshim v0.12.8/go.mod h1:cibQ4BqhJ32FXDwPdQhKhwrwophnh3FuT4nwQZF907w= github.com/anchore/go-struct-converter v0.0.0-20221118182256-c68fdcfa2092 h1:aM1rlcoLz8y5B2r4tTLMiVTrMtpfY0O8EScKJxaSaEc= github.com/anchore/go-struct-converter v0.0.0-20221118182256-c68fdcfa2092/go.mod h1:rYqSE9HbjzpHTI74vwPvae4ZVYZd1lue2ta6xHPdblA= +github.com/aws/aws-sdk-go-v2 v1.41.7 h1:DWpAJt66FmnnaRIOT/8ASTucrvuDPZASqhhLey6tLY8= +github.com/aws/aws-sdk-go-v2 v1.41.7/go.mod h1:4LAfZOPHNVNQEckOACQx60Y8pSRjIkNZQz1w92xpMJc= +github.com/aws/aws-sdk-go-v2/config v1.32.17 h1:FpL4/758/diKwqbytU0prpuiu60fgXKUWCpDJtApclU= +github.com/aws/aws-sdk-go-v2/config v1.32.17/go.mod h1:OXqUMzgXytfoF9JaKkhrOYsyh72t9G+MJH8mMRaexOE= +github.com/aws/aws-sdk-go-v2/credentials v1.19.16 h1:r3RJBuU7X9ibt8RHbMjWE6y60QbKBiII6wSrXnapxSU= +github.com/aws/aws-sdk-go-v2/credentials v1.19.16/go.mod h1:6cx7zqDENJDbBIIWX6P8s0h6hqHC8Avbjh9Dseo27ug= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.23 h1:UuSfcORqNSz/ey3VPRS8TcVH2Ikf0/sC+Hdj400QI6U= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.18.23/go.mod h1:+G/OSGiOFnSOkYloKj/9M35s74LgVAdJBSD5lsFfqKg= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.23 h1:GpT/TrnBYuE5gan2cZbTtvP+JlHsutdmlV2YfEyNde0= +github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.23/go.mod h1:xYWD6BS9ywC5bS3sz9Xh04whO/hzK2plt2Zkyrp4JuA= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.23 h1:bpd8vxhlQi2r1hiueOw02f/duEPTMK59Q4QMAoTTtTo= +github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.23/go.mod h1:15DfR2nw+CRHIk0tqNyifu3G1YdAOy68RftkhMDDwYk= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.24 h1:OQqn11BtaYv1WLUowvcA30MpzIu8Ti4pcLPIIyoKZrA= +github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.24/go.mod h1:X5ZJyfwVrWA96GzPmUCWFQaEARPR7gCrpq2E92PJwAE= +github.com/aws/aws-sdk-go-v2/service/ecr v1.57.2 h1:rHEW02JFJUV2/ttjzyPIvbD0YraqpyU2w6m6DfQUmdg= +github.com/aws/aws-sdk-go-v2/service/ecr v1.57.2/go.mod h1:gNS8pNht4VMzPd4UtQUL3NTUQbjEPLLmb9MqmqrqsCM= +github.com/aws/aws-sdk-go-v2/service/iam v1.53.10 h1:kcN3I3llO7VwIY5w3Pc5FmEonpsr23Ou7Cwk4qf7dik= +github.com/aws/aws-sdk-go-v2/service/iam v1.53.10/go.mod h1:1vkJzjCYC3byO0kIrBqLPzvZpuvYhPXkuyARs6E7tM4= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.9 h1:FLudkZLt5ci0ozzgkVo8BJGwvqNaZbTWb3UcucAateA= +github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.9/go.mod h1:w7wZ/s9qK7c8g4al+UyoF1Sp/Z45UwMGcqIzLWVQHWk= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.23 h1:pbrxO/kuIwgEsOPLkaHu0O+m4fNgLU8B3vxQ+72jTPw= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.23/go.mod h1:/CMNUqoj46HpS3MNRDEDIwcgEnrtZlKRaHNaHxIFpNA= +github.com/aws/aws-sdk-go-v2/service/signin v1.0.11 h1:TdJ+HdzOBhU8+iVAOGUTU63VXopcumCOF1paFulHWZc= +github.com/aws/aws-sdk-go-v2/service/signin v1.0.11/go.mod h1:R82ZRExE/nheo0N+T8zHPcLRTcH8MGsnR3BiVGX0TwI= +github.com/aws/aws-sdk-go-v2/service/sso v1.30.17 h1:7byT8HUWrgoRp6sXjxtZwgOKfhss5fW6SkLBtqzgRoE= +github.com/aws/aws-sdk-go-v2/service/sso v1.30.17/go.mod h1:xNWknVi4Ezm1vg1QsB/5EWpAJURq22uqd38U8qKvOJc= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.35.21 h1:+1Kl1zx6bWi4X7cKi3VYh29h8BvsCoHQEQ6ST9X8w7w= +github.com/aws/aws-sdk-go-v2/service/ssooidc v1.35.21/go.mod h1:4vIRDq+CJB2xFAXZ+YgGUTiEft7oAQlhIs71xcSeuVg= +github.com/aws/aws-sdk-go-v2/service/sts v1.42.1 h1:F/M5Y9I3nwr2IEpshZgh1GeHpOItExNM9L1euNuh/fk= +github.com/aws/aws-sdk-go-v2/service/sts v1.42.1/go.mod h1:mTNxImtovCOEEuD65mKW7DCsL+2gjEH+RPEAexAzAio= +github.com/aws/smithy-go v1.25.1 h1:J8ERsGSU7d+aCmdQur5Txg6bVoYelvQJgtZehD12GkI= +github.com/aws/smithy-go v1.25.1/go.mod h1:YE2RhdIuDbA5E5bTdciG9KrW3+TiEONeUWCqxX9i1Fc= github.com/bahlo/generic-list-go v0.2.0 h1:5sz/EEAK+ls5wF+NeqDpk5+iNdMDXrh3z3nPnH1Wvgk= github.com/bahlo/generic-list-go v0.2.0/go.mod h1:2KvAjgMlE5NNynlg/5iLrrCCZ2+5xWbdbCW3pNTGyYg= github.com/buger/jsonparser v1.1.1 h1:2PnMjfWD7wBILjqQbt530v576A/cAbQvEW9gGIpYMUs= diff --git a/pkg/api/agents.go b/pkg/api/agents.go index 0d5473c..b4f1a1b 100644 --- a/pkg/api/agents.go +++ b/pkg/api/agents.go @@ -19,6 +19,31 @@ import ( "github.com/lyzrai/flow/pkg/storage" ) +// PublicCredential is the scrubbed shape of a stored credential. We +// return non-secret fields (region, accountId, role ARN, kv keys) and +// boolean flags for any sealed values, never the sealed bytes themselves. +type PublicCredential struct { + ID string `json:"id"` + Name string `json:"name"` + Type storage.CredentialType `json:"type"` + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` + + // AWS — all non-secret. + AwsRegion string `json:"awsRegion,omitempty"` + AwsAccountID string `json:"awsAccountId,omitempty"` + AwsCrossAccountRoleArn string `json:"awsCrossAccountRoleArn,omitempty"` + + // GCP — projectId / location are non-secret. HasServiceAccount tells + // the UI whether the SA JSON has been uploaded. + GcpProjectID string `json:"gcpProjectId,omitempty"` + GcpLocation string `json:"gcpLocation,omitempty"` + HasServiceAccount bool `json:"hasServiceAccount,omitempty"` + + // Generic kv — keys are surfaced; values never are. + KvKeys []string `json:"kvKeys,omitempty"` +} + // Agent is the wire shape of an agent record. PAT and webhook secret are // scrubbed; HasPAT and webhook fields surface only the safe parts. type Agent struct { @@ -34,11 +59,43 @@ type Agent struct { AuthStatus storage.AuthStatus `json:"authStatus,omitempty"` AuthCheckedAt *time.Time `json:"authCheckedAt,omitempty"` AttachedPipelines []string `json:"attachedPipelines,omitempty"` - CreatedAt time.Time `json:"createdAt"` - UpdatedAt time.Time `json:"updatedAt"` + Credentials []PublicCredential `json:"credentials,omitempty"` + + CreatedAt time.Time `json:"createdAt"` + UpdatedAt time.Time `json:"updatedAt"` +} + +func publicCredential(c storage.Credential) PublicCredential { + out := PublicCredential{ + ID: c.ID, + Name: c.Name, + Type: c.Type, + CreatedAt: c.CreatedAt, + UpdatedAt: c.UpdatedAt, + } + switch c.Type { + case storage.CredentialAWS: + out.AwsRegion = c.AwsRegion + out.AwsAccountID = c.AwsAccountID + out.AwsCrossAccountRoleArn = c.AwsCrossAccountRoleArn + case storage.CredentialGCP: + out.GcpProjectID = c.GcpProjectID + out.GcpLocation = c.GcpLocation + out.HasServiceAccount = c.GcpServiceAccountSealed != "" + case storage.CredentialKV: + out.KvKeys = make([]string, 0, len(c.KvSealed)) + for k := range c.KvSealed { + out.KvKeys = append(out.KvKeys, k) + } + } + return out } func (s *Server) publicAgent(a *storage.Agent) Agent { + creds := make([]PublicCredential, 0, len(a.Credentials)) + for _, c := range a.Credentials { + creds = append(creds, publicCredential(c)) + } return Agent{ ID: a.ID, Name: a.Name, @@ -52,6 +109,7 @@ func (s *Server) publicAgent(a *storage.Agent) Agent { AuthStatus: a.AuthStatus, AuthCheckedAt: a.AuthCheckedAt, AttachedPipelines: a.AttachedPipelines, + Credentials: creds, CreatedAt: a.CreatedAt, UpdatedAt: a.UpdatedAt, } @@ -391,6 +449,25 @@ func (s *Server) dispatchAgent(ctx context.Context, a *storage.Agent, trigger an } } + // Per-pipeline branch filter — only applies when the trigger came from a + // git push event. Manual triggers fan out to every attached pipeline so + // you can still kick a run from the UI without first hand-editing every + // Trigger node. The filter compares the pushed branch against the + // pipeline's Trigger node `fromBranch` — wildcard ("*", empty) on either + // side disables the filter for that pipeline. + pushedBranch := "" + isPush := false + if items, ok := trigger.([]map[string]any); ok && len(items) > 0 { + if src, _ := items[0]["source"].(string); src == "github_push" { + isPush = true + if b, _ := items[0]["branch"].(string); b != "" { + pushedBranch = b + } else if b, _ := items[0]["ref"].(string); b != "" { + pushedBranch = b + } + } + } + out := make([]string, 0, len(a.AttachedPipelines)) var failures []DispatchFailure for _, pid := range a.AttachedPipelines { @@ -417,6 +494,22 @@ func (s *Server) dispatchAgent(ctx context.Context, a *storage.Agent, trigger an }) continue } + if isPush { + triggerBranch := pipelineTriggerBranch(wf) + if triggerBranch != "" && triggerBranch != "*" && triggerBranch != pushedBranch { + slog.InfoContext(ctx, "agent_dispatch_branch_filtered", + slog.String("pipeline_id", pid), + slog.String("pushed", pushedBranch), + slog.String("trigger_branch", triggerBranch), + ) + failures = append(failures, DispatchFailure{ + PipelineID: pid, + Reason: "branch filtered", + Error: fmt.Sprintf("trigger fromBranch=%q ≠ pushed=%q", triggerBranch, pushedBranch), + }) + continue + } + } runCtx, cancel := context.WithTimeout(ctx, 10*time.Second) execID, err := s.orch.RunAsync(runCtx, &orchestrator.RunRequest{ RequestMeta: orchestrator.RequestMeta{WorkflowID: pid}, @@ -522,16 +615,12 @@ func (s *Server) handleGitHubWebhook(w http.ResponseWriter, r *http.Request) { return } - // Per the agent's configured ref: only dispatch if the branch matches. - if a.Ref != "" && a.Ref != "*" && github.BranchFromRef(push.Ref) != a.Ref { - slog.InfoContext(r.Context(), "github_webhook_ref_mismatch", - slog.String("agent_id", id), - slog.String("event_ref", push.Ref), - slog.String("agent_ref", a.Ref), - ) - writeJSON(w, http.StatusOK, map[string]string{"message": "ref filtered"}) - return - } + // Note: branch filtering happens per-pipeline in dispatchAgent — each + // Trigger node's `fromBranch` decides whether that pipeline runs for + // this push. The agent-level `Ref` is now only the default clone branch + // (used by Build when nothing upstream specifies one), not a webhook + // gate, so chains across branches (Promote main→dev fires the dev + // pipeline) work correctly. if s.orch == nil { writeError(w, http.StatusServiceUnavailable, errors.New("orchestrator not configured")) @@ -575,3 +664,24 @@ func deriveAgentName(raw string) string { } return strings.TrimSpace(raw) } + +// pipelineTriggerBranch returns the `fromBranch` parameter on the workflow's +// Trigger node, or "" if the workflow has no Trigger node or no fromBranch +// configured. Used to filter webhook dispatch so only pipelines matching the +// pushed branch run. If multiple Trigger nodes exist (rare — schema allows it +// but the orchestrator only feeds one), the first wins. +func pipelineTriggerBranch(wf *models.WorkflowDefinition) string { + if wf == nil { + return "" + } + for _, n := range wf.Nodes { + if n.Type != "flow-nodes-base.trigger" { + continue + } + if v, ok := n.Parameters["fromBranch"].(string); ok { + return strings.TrimSpace(v) + } + return "" + } + return "" +} diff --git a/pkg/api/credentials.go b/pkg/api/credentials.go new file mode 100644 index 0000000..ac53fa3 --- /dev/null +++ b/pkg/api/credentials.go @@ -0,0 +1,228 @@ +package api + +import ( + "encoding/json" + "errors" + "fmt" + "net/http" + "strings" + "time" + + "github.com/lyzrai/flow/pkg/secrets" + "github.com/lyzrai/flow/pkg/storage" +) + +// credentialBody is the create/update wire format. Fields are tagged +// per credential type — only the matching ones are read for a given +// `type`. Secrets (gcpServiceAccountJson, kv values) come in as +// plaintext and get sealed before persisting. +type credentialBody struct { + Name string `json:"name"` + Type storage.CredentialType `json:"type"` + + // AWS + AwsRegion string `json:"awsRegion,omitempty"` + AwsAccountID string `json:"awsAccountId,omitempty"` + AwsCrossAccountRoleArn string `json:"awsCrossAccountRoleArn,omitempty"` + + // GCP + GcpProjectID string `json:"gcpProjectId,omitempty"` + GcpLocation string `json:"gcpLocation,omitempty"` + GcpServiceAccountJson string `json:"gcpServiceAccountJson,omitempty"` + + // KV + Kv map[string]string `json:"kv,omitempty"` +} + +// validateAndBuild converts the wire body into a storage.Credential, +// sealing any secret fields. Sets ID + CreatedAt if missing (caller +// pulls those forward on update). +func (b *credentialBody) validateAndBuild() (storage.Credential, error) { + name := strings.TrimSpace(b.Name) + if name == "" { + return storage.Credential{}, errors.New("name is required") + } + now := time.Now().UTC() + c := storage.Credential{ + Name: name, + Type: b.Type, + CreatedAt: now, + UpdatedAt: now, + } + switch b.Type { + case storage.CredentialAWS: + c.AwsRegion = strings.TrimSpace(b.AwsRegion) + c.AwsAccountID = strings.TrimSpace(b.AwsAccountID) + c.AwsCrossAccountRoleArn = strings.TrimSpace(b.AwsCrossAccountRoleArn) + if c.AwsRegion == "" || c.AwsAccountID == "" || c.AwsCrossAccountRoleArn == "" { + return storage.Credential{}, errors.New("aws credential needs region, accountId, and crossAccountRoleArn") + } + case storage.CredentialGCP: + c.GcpProjectID = strings.TrimSpace(b.GcpProjectID) + c.GcpLocation = strings.TrimSpace(b.GcpLocation) + if c.GcpProjectID == "" { + return storage.Credential{}, errors.New("gcp credential needs projectId") + } + if b.GcpServiceAccountJson != "" { + sealed, err := secrets.SealString(b.GcpServiceAccountJson) + if err != nil { + return storage.Credential{}, fmt.Errorf("seal gcp SA: %w", err) + } + c.GcpServiceAccountSealed = sealed + } + case storage.CredentialKV: + if len(b.Kv) == 0 { + return storage.Credential{}, errors.New("kv credential needs at least one key/value") + } + sealed, err := secrets.SealMap(b.Kv) + if err != nil { + return storage.Credential{}, fmt.Errorf("seal kv: %w", err) + } + c.KvSealed = sealed + default: + return storage.Credential{}, fmt.Errorf("unknown credential type %q", b.Type) + } + return c, nil +} + +func (s *Server) handleListCredentials(w http.ResponseWriter, r *http.Request) { + id := r.PathValue("id") + a, err := s.agents.Get(r.Context(), id) + if err != nil { + writeStorageErr(w, err, "agent not found") + return + } + out := make([]PublicCredential, 0, len(a.Credentials)) + for _, c := range a.Credentials { + out = append(out, publicCredential(c)) + } + writeJSON(w, http.StatusOK, out) +} + +func (s *Server) handleCreateCredential(w http.ResponseWriter, r *http.Request) { + if !secrets.IsConfigured() { + writeError(w, http.StatusServiceUnavailable, errors.New("FLOW_SECRET_KEY is not set; refusing to store credentials")) + return + } + id := r.PathValue("id") + a, err := s.agents.Get(r.Context(), id) + if err != nil { + writeStorageErr(w, err, "agent not found") + return + } + + var body credentialBody + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + cred, err := body.validateAndBuild() + if err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + // Reject duplicate names within the agent. + for _, c := range a.Credentials { + if strings.EqualFold(c.Name, cred.Name) { + writeError(w, http.StatusConflict, fmt.Errorf("credential %q already exists; PUT to update", cred.Name)) + return + } + } + cred.ID = newID() + a.Credentials = append(a.Credentials, cred) + a.UpdatedAt = time.Now().UTC() + if err := s.agents.Update(r.Context(), a); err != nil { + writeError(w, http.StatusInternalServerError, err) + return + } + writeJSON(w, http.StatusCreated, publicCredential(cred)) +} + +func (s *Server) handleUpdateCredential(w http.ResponseWriter, r *http.Request) { + if !secrets.IsConfigured() { + writeError(w, http.StatusServiceUnavailable, errors.New("FLOW_SECRET_KEY is not set")) + return + } + id := r.PathValue("id") + name := r.PathValue("name") + a, err := s.agents.Get(r.Context(), id) + if err != nil { + writeStorageErr(w, err, "agent not found") + return + } + + var body credentialBody + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + if body.Name == "" { + body.Name = name + } + updated, err := body.validateAndBuild() + if err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + + idx := -1 + for i, c := range a.Credentials { + if strings.EqualFold(c.Name, name) { + idx = i + break + } + } + if idx == -1 { + writeError(w, http.StatusNotFound, fmt.Errorf("credential %q not found", name)) + return + } + + // Preserve immutable fields; carry forward unchanged secrets when the + // update body left them empty (matches the "edit without re-entering + // the secret" UX). + prev := a.Credentials[idx] + updated.ID = prev.ID + updated.CreatedAt = prev.CreatedAt + if updated.Type == storage.CredentialGCP && updated.GcpServiceAccountSealed == "" { + updated.GcpServiceAccountSealed = prev.GcpServiceAccountSealed + } + if updated.Type == storage.CredentialKV && len(updated.KvSealed) == 0 { + updated.KvSealed = prev.KvSealed + } + + a.Credentials[idx] = updated + a.UpdatedAt = time.Now().UTC() + if err := s.agents.Update(r.Context(), a); err != nil { + writeError(w, http.StatusInternalServerError, err) + return + } + writeJSON(w, http.StatusOK, publicCredential(updated)) +} + +func (s *Server) handleDeleteCredential(w http.ResponseWriter, r *http.Request) { + id := r.PathValue("id") + name := r.PathValue("name") + a, err := s.agents.Get(r.Context(), id) + if err != nil { + writeStorageErr(w, err, "agent not found") + return + } + idx := -1 + for i, c := range a.Credentials { + if strings.EqualFold(c.Name, name) { + idx = i + break + } + } + if idx == -1 { + writeError(w, http.StatusNotFound, fmt.Errorf("credential %q not found", name)) + return + } + a.Credentials = append(a.Credentials[:idx], a.Credentials[idx+1:]...) + a.UpdatedAt = time.Now().UTC() + if err := s.agents.Update(r.Context(), a); err != nil { + writeError(w, http.StatusInternalServerError, err) + return + } + w.WriteHeader(http.StatusNoContent) +} diff --git a/pkg/api/credentials_global.go b/pkg/api/credentials_global.go new file mode 100644 index 0000000..cc127c1 --- /dev/null +++ b/pkg/api/credentials_global.go @@ -0,0 +1,153 @@ +package api + +import ( + "encoding/json" + "errors" + "net/http" + "time" + + "github.com/lyzrai/flow/pkg/secrets" + "github.com/lyzrai/flow/pkg/storage" +) + +// Global (org-wide) credentials. Per-agent overrides live on +// agent.Credentials and are handled by credentials.go; this file owns +// the shared pool. Lookup precedence at runtime (Deploy executor): +// +// node.parameters.credentialName +// → agent.Credentials[name] (agent-specific override) +// → credentials[name] (this global pool) +// +// API surface: +// GET /api/credentials +// POST /api/credentials +// GET /api/credentials/{name} +// PUT /api/credentials/{name} +// DELETE /api/credentials/{name} + +func (s *Server) handleListGlobalCredentials(w http.ResponseWriter, r *http.Request) { + if s.credentials == nil { + writeError(w, http.StatusServiceUnavailable, errors.New("credentials store not configured")) + return + } + list, err := s.credentials.List(r.Context()) + if err != nil { + writeError(w, http.StatusInternalServerError, err) + return + } + out := make([]PublicCredential, 0, len(list)) + for _, c := range list { + out = append(out, publicCredential(*c)) + } + writeJSON(w, http.StatusOK, out) +} + +func (s *Server) handleGetGlobalCredential(w http.ResponseWriter, r *http.Request) { + if s.credentials == nil { + writeError(w, http.StatusServiceUnavailable, errors.New("credentials store not configured")) + return + } + name := r.PathValue("name") + c, err := s.credentials.GetByName(r.Context(), name) + if err != nil { + writeStorageErr(w, err, "credential not found") + return + } + writeJSON(w, http.StatusOK, publicCredential(*c)) +} + +func (s *Server) handleCreateGlobalCredential(w http.ResponseWriter, r *http.Request) { + if !secrets.IsConfigured() { + writeError(w, http.StatusServiceUnavailable, errors.New("FLOW_SECRET_KEY is not set; refusing to store credentials")) + return + } + if s.credentials == nil { + writeError(w, http.StatusServiceUnavailable, errors.New("credentials store not configured")) + return + } + var body credentialBody + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + cred, err := body.validateAndBuild() + if err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + cred.ID = newID() + if err := s.credentials.Create(r.Context(), &cred); err != nil { + if errors.Is(err, storage.ErrAlreadyExists) { + writeError(w, http.StatusConflict, err) + return + } + writeError(w, http.StatusInternalServerError, err) + return + } + writeJSON(w, http.StatusCreated, publicCredential(cred)) +} + +func (s *Server) handleUpdateGlobalCredential(w http.ResponseWriter, r *http.Request) { + if !secrets.IsConfigured() { + writeError(w, http.StatusServiceUnavailable, errors.New("FLOW_SECRET_KEY is not set")) + return + } + if s.credentials == nil { + writeError(w, http.StatusServiceUnavailable, errors.New("credentials store not configured")) + return + } + name := r.PathValue("name") + existing, err := s.credentials.GetByName(r.Context(), name) + if err != nil { + writeStorageErr(w, err, "credential not found") + return + } + + var body credentialBody + if err := json.NewDecoder(r.Body).Decode(&body); err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + if body.Name == "" { + body.Name = name + } + updated, err := body.validateAndBuild() + if err != nil { + writeError(w, http.StatusBadRequest, err) + return + } + // Renaming would require a delete+create against the unique index; + // keep the name immutable through PUT and force the user to recreate + // to rename. Mirrors per-agent behavior. + updated.Name = existing.Name + updated.ID = existing.ID + updated.CreatedAt = existing.CreatedAt + updated.UpdatedAt = time.Now().UTC() + + // Preserve unchanged secrets when the request omits them. + if updated.Type == storage.CredentialGCP && updated.GcpServiceAccountSealed == "" { + updated.GcpServiceAccountSealed = existing.GcpServiceAccountSealed + } + if updated.Type == storage.CredentialKV && len(updated.KvSealed) == 0 { + updated.KvSealed = existing.KvSealed + } + + if err := s.credentials.Update(r.Context(), &updated); err != nil { + writeError(w, http.StatusInternalServerError, err) + return + } + writeJSON(w, http.StatusOK, publicCredential(updated)) +} + +func (s *Server) handleDeleteGlobalCredential(w http.ResponseWriter, r *http.Request) { + if s.credentials == nil { + writeError(w, http.StatusServiceUnavailable, errors.New("credentials store not configured")) + return + } + name := r.PathValue("name") + if err := s.credentials.Delete(r.Context(), name); err != nil { + writeStorageErr(w, err, "credential not found") + return + } + w.WriteHeader(http.StatusNoContent) +} diff --git a/pkg/api/server.go b/pkg/api/server.go index 0cd5e16..e5b11da 100644 --- a/pkg/api/server.go +++ b/pkg/api/server.go @@ -58,9 +58,10 @@ type ServerDeps struct { // "https://abcd.trycloudflare.com"). Used to render webhook callback // URLs that GitHub can hit. Empty means webhook install is disabled. PublicURL string - Pipelines storage.PipelineStore - Runs storage.RunStore - Agents storage.AgentStore + Pipelines storage.PipelineStore + Runs storage.RunStore + Agents storage.AgentStore + Credentials storage.CredentialStore // Events is the in-memory pub/sub bus the orchestrator publishes // per-node lifecycle events to. The SSE handler subscribes per // execution ID. Nil disables /api/executions/{id}/stream. @@ -87,11 +88,12 @@ type Server struct { corsOrigins []string publicURL string - pipelines storage.PipelineStore - runs storage.RunStore - agents storage.AgentStore - events EventSubscriber - logs logstore.Store + pipelines storage.PipelineStore + runs storage.RunStore + agents storage.AgentStore + credentials storage.CredentialStore + events EventSubscriber + logs logstore.Store // runsBus broadcasts run_created events to every UI tab subscribed to // /api/runs/stream. Used so a webhook-triggered run shows up live in @@ -109,10 +111,11 @@ func NewServer(deps ServerDeps) *Server { restateIngres: deps.RestateIngressURL, corsOrigins: deps.CORSOrigins, publicURL: strings.TrimRight(deps.PublicURL, "/"), - pipelines: deps.Pipelines, - runs: deps.Runs, - agents: deps.Agents, - events: deps.Events, + pipelines: deps.Pipelines, + runs: deps.Runs, + agents: deps.Agents, + credentials: deps.Credentials, + events: deps.Events, logs: deps.Logs, runsBus: newRunsBus(), } @@ -167,6 +170,23 @@ func (s *Server) routes() { s.mux.HandleFunc("DELETE /api/agents/{id}/pipelines/{pipelineId}", s.handleDetachPipeline) s.mux.HandleFunc("POST /api/agents/{id}/trigger", s.handleTriggerAgent) + // Credentials API — global org-wide pool. Lookup by name is shared + // across agents/pipelines. Secret fields are AES-GCM sealed at rest + // with FLOW_SECRET_KEY; the API never returns them. + s.mux.HandleFunc("GET /api/credentials", s.handleListGlobalCredentials) + s.mux.HandleFunc("POST /api/credentials", s.handleCreateGlobalCredential) + s.mux.HandleFunc("GET /api/credentials/{name}", s.handleGetGlobalCredential) + s.mux.HandleFunc("PUT /api/credentials/{name}", s.handleUpdateGlobalCredential) + s.mux.HandleFunc("DELETE /api/credentials/{name}", s.handleDeleteGlobalCredential) + + // Per-agent overrides — same shape, but live on the agent doc so a + // pipeline can specialize a credential without touching the global + // pool. + s.mux.HandleFunc("GET /api/agents/{id}/credentials", s.handleListCredentials) + s.mux.HandleFunc("POST /api/agents/{id}/credentials", s.handleCreateCredential) + s.mux.HandleFunc("PUT /api/agents/{id}/credentials/{name}", s.handleUpdateCredential) + s.mux.HandleFunc("DELETE /api/agents/{id}/credentials/{name}", s.handleDeleteCredential) + // Public webhook receiver. GitHub posts here; HMAC signature is the // authentication. Must NOT require CORS / API auth. s.mux.HandleFunc("POST /webhooks/github/{id}", s.handleGitHubWebhook) diff --git a/pkg/awsdeploy/awsdeploy.go b/pkg/awsdeploy/awsdeploy.go new file mode 100644 index 0000000..823f3ef --- /dev/null +++ b/pkg/awsdeploy/awsdeploy.go @@ -0,0 +1,445 @@ +// Package awsdeploy ports the AWS pieces of agent-deploy/build_service: +// cross-account STS AssumeRole, idempotent ECR + AgentCore-runtime-role +// bootstrap, and AgentCore control-plane create-or-update + endpoint +// readiness polling. It is the AWS-target adapter for the Flow Deploy +// node (see pkg/executors/deploy.go). +// +// AWS does not (yet) ship a Go SDK client for bedrock-agentcore-control, +// so AgentCore calls are raw HTTPS signed with SigV4 — the same approach +// the JS reference uses. +package awsdeploy + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "strings" + "time" + + "github.com/aws/aws-sdk-go-v2/aws" + v4 "github.com/aws/aws-sdk-go-v2/aws/signer/v4" + "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/credentials/stscreds" + "github.com/aws/aws-sdk-go-v2/service/ecr" + ecrtypes "github.com/aws/aws-sdk-go-v2/service/ecr/types" + "github.com/aws/aws-sdk-go-v2/service/iam" + "github.com/aws/aws-sdk-go-v2/service/sts" + "github.com/aws/smithy-go" +) + +// Logger is the minimal sink the executor's per-node logger satisfies. +// Pass `engine.NodeLoggerFromContext(ctx)` from the caller. +type Logger interface { + Log(line string) +} + +type nopLogger struct{} + +func (nopLogger) Log(string) {} + +// loggerOr returns lg if non-nil, else a no-op. +func loggerOr(lg Logger) Logger { + if lg == nil { + return nopLogger{} + } + return lg +} + +// Config is everything the deploy executor passes once per call. +type Config struct { + Region string + AccountID string + CrossAccountRoleArn string // role in the customer account; this host assumes it +} + +// AssumeCustomer returns AWS creds for the customer account by calling +// STS AssumeRole on the configured cross-account role. The Flow host's +// own creds (env / IRSA / instance profile) authenticate the AssumeRole. +func AssumeCustomer(ctx context.Context, c Config) (aws.CredentialsProvider, error) { + if c.Region == "" || c.CrossAccountRoleArn == "" { + return nil, errors.New("awsdeploy: region and crossAccountRoleArn are required") + } + hostCfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(c.Region)) + if err != nil { + return nil, fmt.Errorf("load host AWS config: %w", err) + } + stsCli := sts.NewFromConfig(hostCfg) + prov := stscreds.NewAssumeRoleProvider(stsCli, c.CrossAccountRoleArn, func(o *stscreds.AssumeRoleOptions) { + o.RoleSessionName = "flow-deploy" + o.Duration = time.Hour + }) + // Force one resolution up-front so we surface auth errors here, not on + // the first AWS call inside an executor. + if _, err := prov.Retrieve(ctx); err != nil { + return nil, fmt.Errorf("assume role %s: %w", c.CrossAccountRoleArn, err) + } + return aws.NewCredentialsCache(prov), nil +} + +// CustomerConfig is an aws.Config with the cross-account creds attached. +func CustomerConfig(ctx context.Context, c Config, creds aws.CredentialsProvider) (aws.Config, error) { + return config.LoadDefaultConfig(ctx, + config.WithRegion(c.Region), + config.WithCredentialsProvider(creds), + ) +} + +// EnsureECRRepository creates the repo if missing. Idempotent — already- +// exists is treated as success. Returns the repo URI (registry/name). +func EnsureECRRepository(ctx context.Context, awsCfg aws.Config, accountID, region, repoName string, lg Logger) (string, error) { + lg = loggerOr(lg) + cli := ecr.NewFromConfig(awsCfg) + scanOnPush := true + _, err := cli.CreateRepository(ctx, &ecr.CreateRepositoryInput{ + RepositoryName: aws.String(repoName), + ImageScanningConfiguration: &ecrtypes.ImageScanningConfiguration{ + ScanOnPush: scanOnPush, + }, + }) + if err != nil { + var ae smithy.APIError + if errors.As(err, &ae) && ae.ErrorCode() == "RepositoryAlreadyExistsException" { + lg.Log(fmt.Sprintf("[awsdeploy] ECR repo %s already exists", repoName)) + } else { + return "", fmt.Errorf("create ECR repo %s: %w", repoName, err) + } + } else { + lg.Log(fmt.Sprintf("[awsdeploy] ECR repo %s created", repoName)) + } + return fmt.Sprintf("%s.dkr.ecr.%s.amazonaws.com/%s", accountID, region, repoName), nil +} + +// EnsureAgentCoreRuntimeRole creates the shared `agentcore-runtime-role` +// (per the JS reference) if missing, and attaches a runtime policy. The +// role's trust policy allows bedrock-agentcore.amazonaws.com to assume it, +// scoped to the customer account. Idempotent. +func EnsureAgentCoreRuntimeRole(ctx context.Context, awsCfg aws.Config, accountID, region string, lg Logger) (string, error) { + lg = loggerOr(lg) + const roleName = "agentcore-runtime-role" + const policyName = "AgentCoreRuntimePolicy" + + cli := iam.NewFromConfig(awsCfg) + if out, err := cli.GetRole(ctx, &iam.GetRoleInput{RoleName: aws.String(roleName)}); err == nil { + lg.Log(fmt.Sprintf("[awsdeploy] IAM role %s exists", roleName)) + return aws.ToString(out.Role.Arn), nil + } else { + var ae smithy.APIError + if !errors.As(err, &ae) || ae.ErrorCode() != "NoSuchEntity" { + return "", fmt.Errorf("get role: %w", err) + } + } + + trust := mustJSON(map[string]any{ + "Version": "2012-10-17", + "Statement": []map[string]any{{ + "Effect": "Allow", + "Principal": map[string]any{"Service": "bedrock-agentcore.amazonaws.com"}, + "Action": "sts:AssumeRole", + "Condition": map[string]any{"StringEquals": map[string]string{"aws:SourceAccount": accountID}}, + }}, + }) + createOut, err := cli.CreateRole(ctx, &iam.CreateRoleInput{ + RoleName: aws.String(roleName), + AssumeRolePolicyDocument: aws.String(trust), + Description: aws.String("Shared AgentCore runtime role created by Flow"), + }) + if err != nil { + var ae smithy.APIError + if errors.As(err, &ae) && ae.ErrorCode() == "EntityAlreadyExists" { + // Race: another deploy created it between Get and Create. + lg.Log(fmt.Sprintf("[awsdeploy] IAM role %s won the race", roleName)) + out, gerr := cli.GetRole(ctx, &iam.GetRoleInput{RoleName: aws.String(roleName)}) + if gerr != nil { + return "", gerr + } + return aws.ToString(out.Role.Arn), nil + } + return "", fmt.Errorf("create role: %w", err) + } + roleArn := aws.ToString(createOut.Role.Arn) + + policy := mustJSON(map[string]any{ + "Version": "2012-10-17", + "Statement": []map[string]any{ + {"Effect": "Allow", "Action": []string{"ecr:BatchGetImage", "ecr:GetDownloadUrlForLayer", "ecr:GetAuthorizationToken"}, "Resource": "*"}, + {"Effect": "Allow", "Action": []string{"logs:CreateLogGroup", "logs:CreateLogStream", "logs:PutLogEvents", "logs:DescribeLogStreams", "logs:DescribeLogGroups"}, "Resource": "*"}, + {"Effect": "Allow", "Action": []string{"xray:PutTraceSegments", "xray:PutTelemetryRecords", "xray:GetSamplingRules", "xray:GetSamplingTargets"}, "Resource": "*"}, + {"Effect": "Allow", "Action": "cloudwatch:PutMetricData", "Resource": "*", "Condition": map[string]any{"StringEquals": map[string]string{"cloudwatch:namespace": "bedrock-agentcore"}}}, + {"Effect": "Allow", "Action": []string{"bedrock:InvokeModel", "bedrock:InvokeModelWithResponseStream"}, "Resource": fmt.Sprintf("arn:aws:bedrock:%s::foundation-model/*", region)}, + {"Effect": "Allow", "Action": []string{"bedrock-agentcore:GetWorkloadAccessToken", "bedrock-agentcore:GetWorkloadAccessTokenForJWT", "bedrock-agentcore:GetWorkloadAccessTokenForUserId"}, "Resource": "*"}, + }, + }) + if _, err := cli.PutRolePolicy(ctx, &iam.PutRolePolicyInput{ + RoleName: aws.String(roleName), + PolicyName: aws.String(policyName), + PolicyDocument: aws.String(policy), + }); err != nil { + return "", fmt.Errorf("put role policy: %w", err) + } + lg.Log(fmt.Sprintf("[awsdeploy] IAM role %s created; sleeping 10s for propagation", roleName)) + select { + case <-ctx.Done(): + return "", ctx.Err() + case <-time.After(10 * time.Second): + } + return roleArn, nil +} + +// AgentRuntimeResult holds the post-deploy identifiers returned to the +// Deploy executor's __deploy summary. EndpointArn + InvokeURL are what +// users actually want at the end of the run. +type AgentRuntimeResult struct { + AgentRuntimeID string + AgentRuntimeArn string + EndpointArn string + InvokeURL string +} + +// CreateOrUpdateAgentRuntime puts a new AgentCore runtime (or updates the +// existing one with the same name) and waits for the DEFAULT endpoint to +// reach READY. Mirrors deploy_to_agentcore.js exactly. timeout caps the +// readiness wait; total time also bounded by ctx. +func CreateOrUpdateAgentRuntime( + ctx context.Context, + c Config, + creds aws.CredentialsProvider, + runtimeName, image, runtimeRoleArn string, + envVars map[string]string, + timeout time.Duration, + lg Logger, +) (AgentRuntimeResult, error) { + lg = loggerOr(lg) + if c.Region == "" { + return AgentRuntimeResult{}, errors.New("region required") + } + if image == "" { + return AgentRuntimeResult{}, errors.New("image required") + } + if runtimeRoleArn == "" { + return AgentRuntimeResult{}, errors.New("runtimeRoleArn required") + } + + // AgentCore enforces [a-zA-Z0-9_]{1,48} + agentName := sanitizeAgentName(runtimeName) + createBody := map[string]any{ + "agentRuntimeName": agentName, + "agentRuntimeArtifact": map[string]any{ + "containerConfiguration": map[string]any{"containerUri": image}, + }, + "roleArn": runtimeRoleArn, + "networkConfiguration": map[string]any{"networkMode": "PUBLIC"}, + } + if len(envVars) > 0 { + createBody["environmentVariables"] = envVars + } + + var ( + agentID string + agentArn string + ) + + // Try create. + respCreate, status, err := agentcoreRequest(ctx, "PUT", "/runtimes/", createBody, c.Region, creds) + switch { + case err == nil: + agentID, _ = respCreate["agentRuntimeId"].(string) + agentArn, _ = respCreate["agentRuntimeArn"].(string) + lg.Log(fmt.Sprintf("[awsdeploy] AgentCore created runtime %s", agentID)) + case status == http.StatusConflict || isConflict(err): + // Update path: list, find by name, update. + lg.Log(fmt.Sprintf("[awsdeploy] AgentCore runtime %q exists; updating", agentName)) + listResp, _, lerr := agentcoreRequest(ctx, "POST", "/runtimes/", map[string]any{"maxResults": 100}, c.Region, creds) + if lerr != nil { + return AgentRuntimeResult{}, fmt.Errorf("list runtimes: %w", lerr) + } + runtimes, _ := listResp["agentRuntimes"].([]any) + for _, r := range runtimes { + rm, ok := r.(map[string]any) + if !ok { + continue + } + if rm["agentRuntimeName"] == agentName { + agentID, _ = rm["agentRuntimeId"].(string) + agentArn, _ = rm["agentRuntimeArn"].(string) + break + } + } + if agentID == "" { + return AgentRuntimeResult{}, fmt.Errorf("conflict but runtime %q not found in list", agentName) + } + updateBody := map[string]any{ + "agentRuntimeArtifact": map[string]any{ + "containerConfiguration": map[string]any{"containerUri": image}, + }, + "roleArn": runtimeRoleArn, + "networkConfiguration": map[string]any{"networkMode": "PUBLIC"}, + } + if len(envVars) > 0 { + updateBody["environmentVariables"] = envVars + } + if _, _, uerr := agentcoreRequest(ctx, "PUT", "/runtimes/"+agentID+"/", updateBody, c.Region, creds); uerr != nil { + return AgentRuntimeResult{}, fmt.Errorf("update runtime: %w", uerr) + } + lg.Log(fmt.Sprintf("[awsdeploy] AgentCore updated runtime %s", agentID)) + default: + return AgentRuntimeResult{}, fmt.Errorf("create runtime: %w", err) + } + + // Wait for DEFAULT endpoint READY. + endpointArn, werr := waitEndpointReady(ctx, agentID, c.Region, creds, timeout, lg) + if werr != nil { + return AgentRuntimeResult{ + AgentRuntimeID: agentID, + AgentRuntimeArn: agentArn, + }, werr + } + invokeURL := fmt.Sprintf("https://bedrock-agentcore.%s.amazonaws.com/runtimes/%s/invocations", c.Region, agentID) + lg.Log(fmt.Sprintf("[awsdeploy] AgentCore endpoint READY: %s", invokeURL)) + return AgentRuntimeResult{ + AgentRuntimeID: agentID, + AgentRuntimeArn: agentArn, + EndpointArn: endpointArn, + InvokeURL: invokeURL, + }, nil +} + +func waitEndpointReady(ctx context.Context, agentID, region string, creds aws.CredentialsProvider, timeout time.Duration, lg Logger) (string, error) { + if timeout <= 0 { + timeout = 5 * time.Minute + } + deadline := time.Now().Add(timeout) + path := fmt.Sprintf("/runtimes/%s/runtime-endpoints/DEFAULT/", agentID) + for { + if time.Now().After(deadline) { + return "", fmt.Errorf("endpoint not READY within %s", timeout) + } + resp, status, err := agentcoreRequest(ctx, "GET", path, nil, region, creds) + if err != nil { + if status == http.StatusNotFound { + lg.Log("[awsdeploy] endpoint not provisioned yet, waiting...") + } else { + return "", err + } + } else { + s, _ := resp["status"].(string) + lg.Log(fmt.Sprintf("[awsdeploy] endpoint status: %s", s)) + if s == "READY" { + arn, _ := resp["agentRuntimeEndpointArn"].(string) + return arn, nil + } + if strings.Contains(s, "FAILED") { + reason, _ := resp["failureReason"].(string) + return "", fmt.Errorf("endpoint failed: %s (%s)", s, reason) + } + } + select { + case <-ctx.Done(): + return "", ctx.Err() + case <-time.After(5 * time.Second): + } + } +} + +// agentcoreRequest signs and sends a request to bedrock-agentcore-control. +// Returns the parsed JSON response, HTTP status, and a non-nil error if +// status >= 400 or transport failed. +func agentcoreRequest( + ctx context.Context, + method, path string, + body map[string]any, + region string, + creds aws.CredentialsProvider, +) (map[string]any, int, error) { + host := fmt.Sprintf("bedrock-agentcore-control.%s.amazonaws.com", region) + endpoint := url.URL{Scheme: "https", Host: host, Path: path} + + var bodyBytes []byte + if body != nil { + var err error + bodyBytes, err = json.Marshal(body) + if err != nil { + return nil, 0, err + } + } + req, err := http.NewRequestWithContext(ctx, method, endpoint.String(), bytes.NewReader(bodyBytes)) + if err != nil { + return nil, 0, err + } + req.Header.Set("content-type", "application/json") + req.Host = host + + c, err := creds.Retrieve(ctx) + if err != nil { + return nil, 0, fmt.Errorf("retrieve creds: %w", err) + } + hash := sha256Hex(bodyBytes) + signer := v4.NewSigner() + if err := signer.SignHTTP(ctx, c, req, hash, "bedrock-agentcore", region, time.Now()); err != nil { + return nil, 0, fmt.Errorf("sigv4 sign: %w", err) + } + + resp, err := http.DefaultClient.Do(req) + if err != nil { + return nil, 0, err + } + defer resp.Body.Close() + raw, _ := io.ReadAll(resp.Body) + + var parsed map[string]any + if len(raw) > 0 { + _ = json.Unmarshal(raw, &parsed) + } + if resp.StatusCode >= 400 { + msg := "" + if parsed != nil { + if m, ok := parsed["message"].(string); ok { + msg = m + } else if m, ok := parsed["Message"].(string); ok { + msg = m + } + } + if msg == "" { + msg = string(raw) + } + return parsed, resp.StatusCode, fmt.Errorf("agentcore HTTP %d: %s", resp.StatusCode, msg) + } + return parsed, resp.StatusCode, nil +} + +func isConflict(err error) bool { + return err != nil && strings.Contains(strings.ToLower(err.Error()), "conflict") +} + +func sha256Hex(b []byte) string { + h := sha256.Sum256(b) + return hex.EncodeToString(h[:]) +} + +func sanitizeAgentName(in string) string { + var b strings.Builder + for _, r := range in { + switch { + case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9', r == '_': + b.WriteRune(r) + default: + b.WriteRune('_') + } + } + out := b.String() + if len(out) > 48 { + out = out[:48] + } + return out +} + +func mustJSON(v any) string { + b, _ := json.Marshal(v) + return string(b) +} diff --git a/pkg/executors/approval.go b/pkg/executors/approval.go index 86a946e..397a5a9 100644 --- a/pkg/executors/approval.go +++ b/pkg/executors/approval.go @@ -56,16 +56,19 @@ func (e *ApprovalExecutor) Execute( restate.Set(rctx, "pending_approval_node", node.Name) restate.Set(rctx, "pending_approval_id", awakeableID) - // Approval context for the reviewer UI: resolved message + input data. + // Approval context for the reviewer UI: resolved message + reason + + // input data. approvalCtx := map[string]any{} + resolved := node.Parameters if execCtx != nil { - resolved := engine.ResolveExpressions(node.Parameters, execCtx, node.Name) - if msg, ok := resolved["message"].(string); ok && msg != "" { - approvalCtx["message"] = msg - } - } else if msg, ok := node.Parameters["message"].(string); ok && msg != "" { + resolved = engine.ResolveExpressions(node.Parameters, execCtx, node.Name) + } + if msg, ok := resolved["message"].(string); ok && msg != "" { approvalCtx["message"] = msg } + if reason, ok := resolved["reason"].(string); ok && reason != "" { + approvalCtx["reason"] = reason + } if len(inputItems) == 1 { approvalCtx["inputs"] = map[string]any(inputItems[0]) } else if len(inputItems) > 1 { diff --git a/pkg/executors/build.go b/pkg/executors/build.go index 00da397..7a5ed40 100644 --- a/pkg/executors/build.go +++ b/pkg/executors/build.go @@ -80,7 +80,14 @@ func (e *BuildExecutor) Execute(ctx context.Context, node models.NodeDef, inputs // payload (or older trigger record) carried "refs/heads/main", trim it // so the clone doesn't fail with "Remote branch refs/heads/main not // found in upstream origin". - ref := stripRefsHeads(strFirst(strFromAny(trigger["ref"]), a.Ref, "main")) + // Prefer the Trigger node's explicit `fromBranch` over the older + // webhook `ref` field, then fall back to the agent's configured ref. + ref := stripRefsHeads(strFirst( + strFromAny(trigger["fromBranch"]), + strFromAny(trigger["ref"]), + a.Ref, + "main", + )) cloneDir, cleanup, err := cloneRepo(ctx, a, ref, commitSHA, time.Duration(timeoutSec)*time.Second) if err != nil { diff --git a/pkg/executors/deploy.go b/pkg/executors/deploy.go new file mode 100644 index 0000000..a5e7486 --- /dev/null +++ b/pkg/executors/deploy.go @@ -0,0 +1,272 @@ +package executors + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "github.com/lyzrai/flow/pkg/awsdeploy" + "github.com/lyzrai/flow/pkg/engine" + "github.com/lyzrai/flow/pkg/models" + "github.com/lyzrai/flow/pkg/storage" +) + +// DeployExecutor deploys an already-built container image to a managed +// agent runtime. Today only the AWS Bedrock AgentCore target is wired +// (target=agentcore); k8s and Vertex come later (the form / catalog show +// them as "coming soon"). +// +// On success, every output item gets a `__deploy` summary that includes +// the invokeUrl — the public HTTPS endpoint clients hit to talk to the +// deployed agent. A final node-log line restates that URL so it is +// visible without expanding the run output. +// +// Image-source resolution order: +// 1. node.parameters.image (explicit override) +// 2. upstream `__push.copies[0].imageRef` (Push node — the canonical case) +// 3. upstream `__push.dst` (legacy single-target push shape) +// 4. upstream `__build.image` (Build node, no Push step) +// +// AWS creds: the agent record must carry awsRegion + awsAccountId + +// awsCrossAccountRoleArn. The Flow host's own AWS identity (env / IRSA / +// instance profile) calls STS:AssumeRole on that cross-account role; the +// resulting temp creds drive ECR / IAM / AgentCore. +type DeployExecutor struct { + Agents storage.AgentStore + Credentials storage.CredentialStore // optional — global credential pool fallback +} + +func (e *DeployExecutor) Execute(ctx context.Context, node models.NodeDef, inputs [][]models.Item, _ *engine.ExecutionContext) (map[int][]models.Item, error) { + if e.Agents == nil { + return nil, errors.New("deploy: AgentStore not configured") + } + logger := engine.NodeLoggerFromContext(ctx) + + target := strings.ToLower(strParam(node.Parameters, "target", "agentcore")) + switch target { + case "agentcore": + // supported below + case "k8s", "kubernetes", "vertex", "": + return nil, fmt.Errorf("deploy: target %q not yet implemented (only \"agentcore\" today)", target) + default: + return nil, fmt.Errorf("deploy: unknown target %q", target) + } + + trigger := firstItem(inputs) + agentID, _ := trigger["agentId"].(string) + if agentID == "" { + return nil, errors.New("deploy: trigger payload missing agentId") + } + a, err := e.Agents.Get(ctx, agentID) + if err != nil { + return nil, fmt.Errorf("deploy: load agent %q: %w", agentID, err) + } + credName := strParam(node.Parameters, "credentialName", "aws") + cred, credScope, err := e.lookupCredential(ctx, a, credName) + if err != nil { + return nil, fmt.Errorf("deploy: %w", err) + } + logger.Log(fmt.Sprintf("[deploy] using %s credential %q", credScope, credName)) + if cred.Type != storage.CredentialAWS { + return nil, fmt.Errorf("deploy: credential %q is type %q; target=agentcore needs an aws credential", credName, cred.Type) + } + if cred.AwsRegion == "" || cred.AwsAccountID == "" || cred.AwsCrossAccountRoleArn == "" { + return nil, fmt.Errorf("deploy: aws credential %q is missing region / accountId / crossAccountRoleArn", credName) + } + + image := resolveDeployImage(node.Parameters, inputs) + if image == "" { + return nil, errors.New("deploy: no image to deploy — Push must run upstream, or set `image` on the node") + } + + runtimeName := strParam(node.Parameters, "runtimeName", "") + if runtimeName == "" { + runtimeName = deriveRuntimeName(a) + } + envVars := resolveEnvVars(node.Parameters) + timeoutSec := intParam(node.Parameters, "timeoutSeconds", 600) + hardCtx, cancel := context.WithTimeout(ctx, time.Duration(timeoutSec)*time.Second) + defer cancel() + + logger.Log(fmt.Sprintf("[deploy:agentcore] credential=%s account=%s region=%s runtime=%s image=%s", + cred.Name, cred.AwsAccountID, cred.AwsRegion, runtimeName, image)) + + awsCfg := awsdeploy.Config{ + Region: cred.AwsRegion, + AccountID: cred.AwsAccountID, + CrossAccountRoleArn: cred.AwsCrossAccountRoleArn, + } + creds, err := awsdeploy.AssumeCustomer(hardCtx, awsCfg) + if err != nil { + return nil, fmt.Errorf("deploy: %w", err) + } + customerCfg, err := awsdeploy.CustomerConfig(hardCtx, awsCfg, creds) + if err != nil { + return nil, fmt.Errorf("deploy: build customer aws config: %w", err) + } + + // Bootstrap (idempotent): ECR repo + shared agentcore-runtime-role. + // Repo name follows the JS reference: lowercased "agentcore-". + ecrRepoName := strings.ToLower("agentcore-" + runtimeName) + if _, err := awsdeploy.EnsureECRRepository(hardCtx, customerCfg, cred.AwsAccountID, cred.AwsRegion, ecrRepoName, logger); err != nil { + return nil, fmt.Errorf("deploy: %w", err) + } + runtimeRoleArn, err := awsdeploy.EnsureAgentCoreRuntimeRole(hardCtx, customerCfg, cred.AwsAccountID, cred.AwsRegion, logger) + if err != nil { + return nil, fmt.Errorf("deploy: %w", err) + } + + // Create or update the AgentCore runtime + wait for endpoint READY. + res, err := awsdeploy.CreateOrUpdateAgentRuntime( + hardCtx, + awsCfg, + creds, + runtimeName, + image, + runtimeRoleArn, + envVars, + time.Duration(timeoutSec)*time.Second, + logger, + ) + if err != nil { + return nil, fmt.Errorf("deploy: %w", err) + } + + // Make the endpoint impossible to miss in the run timeline. + logger.Log("══════════════════════════════════════════════════════════════════") + logger.Log(fmt.Sprintf("✅ Deploy complete. Invoke URL: %s", res.InvokeURL)) + logger.Log(fmt.Sprintf(" Agent runtime: %s", res.AgentRuntimeArn)) + logger.Log(fmt.Sprintf(" Endpoint ARN : %s", res.EndpointArn)) + logger.Log("══════════════════════════════════════════════════════════════════") + + summary := map[string]any{ + "target": "agentcore", + "agentId": agentID, + "agentName": a.Name, + "credentialName": cred.Name, + "credentialScope": credScope, + "region": cred.AwsRegion, + "accountId": cred.AwsAccountID, + "runtimeName": runtimeName, + "image": image, + "agentRuntimeId": res.AgentRuntimeID, + "agentRuntimeArn": res.AgentRuntimeArn, + "endpointArn": res.EndpointArn, + "invokeUrl": res.InvokeURL, + "runtimeRoleArn": runtimeRoleArn, + "finished_at": time.Now().UTC(), + } + + out := make([]models.Item, 0) + for _, in := range inputs { + for _, it := range in { + ci := copyItem(it) + ci["__deploy"] = summary + out = append(out, ci) + } + } + if len(out) == 0 { + out = append(out, models.Item{"__deploy": summary}) + } + return map[int][]models.Item{0: out}, nil +} + +// resolveDeployImage walks the upstream items in priority order (Push > +// Build) and picks the first image ref it finds. Node-level `image` +// param wins over both. Returns "" if nothing is set. +func resolveDeployImage(params map[string]any, inputs [][]models.Item) string { + if v := strParam(params, "image", ""); v != "" { + return v + } + for _, in := range inputs { + for _, it := range in { + // __push.copies[0].imageRef + if push, ok := it["__push"].(map[string]any); ok { + if copies, ok := push["copies"].([]any); ok && len(copies) > 0 { + if first, ok := copies[0].(map[string]any); ok { + if ref, ok := first["imageRef"].(string); ok && ref != "" { + return ref + } + } + } + // Legacy single-target shape: __push.dst + if dst, ok := push["dst"].(string); ok && dst != "" { + return dst + } + } + // __build.image + if build, ok := it["__build"].(map[string]any); ok { + if ref, ok := build["image"].(string); ok && ref != "" { + return ref + } + } + } + } + return "" +} + +// resolveEnvVars accepts either a map[string]any (UI form's natural shape) +// or a JSON string. Empty values are dropped — AgentCore doesn't accept +// empty-string values on environmentVariables. +func resolveEnvVars(params map[string]any) map[string]string { + out := map[string]string{} + raw, ok := params["envVars"] + if !ok { + return out + } + if m, ok := raw.(map[string]any); ok { + for k, v := range m { + if s, ok := v.(string); ok && s != "" { + out[k] = s + } + } + } + return out +} + +// deriveRuntimeName: agent.Name is "owner/repo" — AgentCore enforces +// [a-zA-Z0-9_]{1,48}, so we replace `/` with `_` and let the adapter's +// sanitizer handle anything else. Empty agent name falls back to the ID. +func deriveRuntimeName(a *storage.Agent) string { + n := a.Name + if n == "" { + n = a.ID + } + return strings.ReplaceAll(n, "/", "_") +} + +// lookupCredential resolves a credential by name with the documented +// precedence: agent-level override first, then the global pool. Returns +// the resolved credential, a scope label ("agent" | "global") for the +// __deploy summary, and a typed error. +// +// We require the credential to be of type AWS for now since that's the +// only target the executor implements; tightening here gives a clearer +// error than letting empty fields blow up further down. +func (e *DeployExecutor) lookupCredential(ctx context.Context, a *storage.Agent, name string) (storage.Credential, string, error) { + if c, ok := a.LookupCredential(name); ok { + if c.Type != storage.CredentialAWS { + return storage.Credential{}, "", fmt.Errorf("agent credential %q is type %q; target=agentcore needs an aws credential", name, c.Type) + } + if c.AwsRegion == "" || c.AwsAccountID == "" || c.AwsCrossAccountRoleArn == "" { + return storage.Credential{}, "", fmt.Errorf("agent aws credential %q is missing region / accountId / crossAccountRoleArn", name) + } + return c, "agent", nil + } + if e.Credentials == nil { + return storage.Credential{}, "", fmt.Errorf("no credential %q on agent and global credential store not configured", name) + } + gc, err := e.Credentials.GetByName(ctx, name) + if err != nil { + return storage.Credential{}, "", fmt.Errorf("no credential %q on agent or in the global pool: %w", name, err) + } + if gc.Type != storage.CredentialAWS { + return storage.Credential{}, "", fmt.Errorf("global credential %q is type %q; target=agentcore needs an aws credential", name, gc.Type) + } + if gc.AwsRegion == "" || gc.AwsAccountID == "" || gc.AwsCrossAccountRoleArn == "" { + return storage.Credential{}, "", fmt.Errorf("global aws credential %q is missing region / accountId / crossAccountRoleArn", name) + } + return *gc, "global", nil +} diff --git a/pkg/executors/promote.go b/pkg/executors/promote.go new file mode 100644 index 0000000..52c56c3 --- /dev/null +++ b/pkg/executors/promote.go @@ -0,0 +1,190 @@ +package executors + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "github.com/lyzrai/flow/pkg/engine" + "github.com/lyzrai/flow/pkg/github" + "github.com/lyzrai/flow/pkg/models" + "github.com/lyzrai/flow/pkg/storage" +) + +// PromoteExecutor promotes the agent's source code by either opening a +// pull request from `fromBranch` → `toBranch`, or merging the two +// branches directly via the GitHub API. This is GitOps promotion, not +// image-tag promotion — promotion is recorded as a real commit / PR in +// the agent's repo so the audit trail lives where reviewers expect it. +// +// Source of branches (precedence): +// 1. node.parameters.fromBranch / toBranch +// 2. trigger payload `fromBranch` / `toBranch` (Trigger node carries these) +// 3. trigger payload `ref` for fromBranch +// 4. agent.Ref for fromBranch +// 5. "main" for fromBranch, error if no toBranch +// +// Modes: +// - open-pr (default): POST /repos/.../pulls; idempotent (re-finds an +// existing open PR for the same head/base pair). +// - merge : POST /repos/.../merges → fast-forward / merge-commit. +// - merge-pr : open the PR if needed, then PUT /pulls/{n}/merge. +// +// Parameters: +// - mode "open-pr" | "merge" | "merge-pr" (default open-pr) +// - fromBranch string override +// - toBranch string override +// - title string PR title (open-pr / merge-pr) +// - body string PR body (open-pr / merge-pr) +// - mergeMethod "merge" | "squash" | "rebase" (merge-pr only) +// - timeoutSeconds number (default 60) +type PromoteExecutor struct { + Agents storage.AgentStore +} + +func (e *PromoteExecutor) Execute(ctx context.Context, node models.NodeDef, inputs [][]models.Item, _ *engine.ExecutionContext) (map[int][]models.Item, error) { + if e.Agents == nil { + return nil, errors.New("promote: AgentStore not configured") + } + logger := engine.NodeLoggerFromContext(ctx) + + trigger := firstItem(inputs) + agentID, _ := trigger["agentId"].(string) + if agentID == "" { + return nil, errors.New("promote: trigger payload missing agentId") + } + a, err := e.Agents.Get(ctx, agentID) + if err != nil { + return nil, fmt.Errorf("promote: load agent %q: %w", agentID, err) + } + repo, err := github.ParseRepo(a.RepoURL) + if err != nil { + return nil, fmt.Errorf("promote: parse repo %q: %w", a.RepoURL, err) + } + if a.PAT == "" { + return nil, errors.New("promote: agent has no PAT (need 'repo' scope to open PRs / merge)") + } + + mode := strings.ToLower(strParam(node.Parameters, "mode", "open-pr")) + from := stripRefsHeads(resolveBranch( + strParam(node.Parameters, "fromBranch", ""), + strFromAny(trigger["fromBranch"]), + strFromAny(trigger["ref"]), + a.Ref, + "main", + )) + to := stripRefsHeads(resolveBranch( + strParam(node.Parameters, "toBranch", ""), + strFromAny(trigger["toBranch"]), + )) + if to == "" { + return nil, errors.New("promote: no toBranch (set on Trigger node or Promote node)") + } + if from == to { + return nil, fmt.Errorf("promote: fromBranch and toBranch are both %q", from) + } + + timeoutSec := intParam(node.Parameters, "timeoutSeconds", 60) + hardCtx, cancel := context.WithTimeout(ctx, time.Duration(timeoutSec)*time.Second) + defer cancel() + + logger.Log(fmt.Sprintf("[promote:%s] %s/%s: %s → %s", + mode, repo.Owner, repo.Name, from, to)) + + cli := github.NewClient(a.PAT) + title := strParam(node.Parameters, "title", "") + body := strParam(node.Parameters, "body", "") + commitMsg := firstNonEmptyStr(title, fmt.Sprintf("Promote %s → %s", from, to)) + + summary := map[string]any{ + "mode": mode, + "agentId": agentID, + "agentName": a.Name, + "repo": repo.String(), + "fromBranch": from, + "toBranch": to, + "finished_at": time.Now().UTC(), + } + + switch mode { + case "open-pr", "": + pr, err := cli.OpenPullRequest(hardCtx, repo, from, to, title, body) + if err != nil { + return nil, fmt.Errorf("promote: %w", err) + } + logger.Log(fmt.Sprintf("[promote] PR #%d opened: %s", pr.Number, pr.HTMLURL)) + summary["prNumber"] = pr.Number + summary["prUrl"] = pr.HTMLURL + summary["prState"] = pr.State + + case "merge": + sha, err := cli.MergeBranches(hardCtx, repo, to, from, commitMsg) + if err != nil { + return nil, fmt.Errorf("promote: %w", err) + } + if sha == "" { + logger.Log("[promote] branches already up-to-date; nothing to merge") + summary["upToDate"] = true + } else { + logger.Log(fmt.Sprintf("[promote] merged %s → %s as %s", from, to, sha)) + summary["mergeSha"] = sha + } + + case "merge-pr": + pr, err := cli.OpenPullRequest(hardCtx, repo, from, to, title, body) + if err != nil { + return nil, fmt.Errorf("promote (open): %w", err) + } + summary["prNumber"] = pr.Number + summary["prUrl"] = pr.HTMLURL + + method := strings.ToLower(strParam(node.Parameters, "mergeMethod", "merge")) + sha, err := cli.MergePullRequest(hardCtx, repo, pr.Number, commitMsg, method) + if err != nil { + return nil, fmt.Errorf("promote (merge PR #%d): %w", pr.Number, err) + } + logger.Log(fmt.Sprintf("[promote] PR #%d merged as %s", pr.Number, sha)) + summary["mergeSha"] = sha + summary["mergeMethod"] = method + + default: + return nil, fmt.Errorf("promote: unknown mode %q (want open-pr|merge|merge-pr)", mode) + } + + out := make([]models.Item, 0) + for _, in := range inputs { + for _, it := range in { + ci := copyItem(it) + ci["__promote"] = summary + out = append(out, ci) + } + } + if len(out) == 0 { + out = append(out, models.Item{"__promote": summary}) + } + return map[int][]models.Item{0: out}, nil +} + +// resolveBranch returns the first non-empty trimmed string in the given +// list. Empty by default — the caller decides whether that's an error. +func resolveBranch(candidates ...string) string { + for _, c := range candidates { + if s := strings.TrimSpace(c); s != "" { + return s + } + } + return "" +} + +// firstNonEmptyStr is the strParam-friendly variant of strFirst that takes +// any number of strings (kept local so we don't conflict with Push's helper). +func firstNonEmptyStr(ss ...string) string { + for _, s := range ss { + if strings.TrimSpace(s) != "" { + return s + } + } + return "" +} diff --git a/pkg/executors/registry.go b/pkg/executors/registry.go index 2e5003c..3ccd05d 100644 --- a/pkg/executors/registry.go +++ b/pkg/executors/registry.go @@ -35,7 +35,8 @@ func Get(nodeType string) (NodeExecutor, error) { // RegistryDeps holds optional dependencies for executors that need external access. type RegistryDeps struct { WorkflowLoader WorkflowLoaderFunc - Agents storage.AgentStore // used by the Build executor to look up the source repo + PAT + Agents storage.AgentStore // used by Build (clone) and Deploy (lookup agent-scoped creds) + Credentials storage.CredentialStore // global credential pool — Deploy falls back to this } // WorkflowLoaderFunc loads a workflow definition by ID from storage. @@ -64,8 +65,8 @@ func RegisterAll(deps ...RegistryDeps) { Register("flow-nodes-base.test", &TestExecutor{}) Register("flow-nodes-base.eval", &EvalExecutor{}) Register("flow-nodes-base.policy", &PolicyExecutor{}) - Register("flow-nodes-base.deploy", &DeployExecutor{}) - Register("flow-nodes-base.promote", &PromoteExecutor{}) + Register("flow-nodes-base.deploy", &DeployExecutor{Agents: d.Agents, Credentials: d.Credentials}) + Register("flow-nodes-base.promote", &PromoteExecutor{Agents: d.Agents}) Register("flow-nodes-base.rollback", &RollbackExecutor{}) } diff --git a/pkg/executors/stubs.go b/pkg/executors/stubs.go index b976954..69eedae 100644 --- a/pkg/executors/stubs.go +++ b/pkg/executors/stubs.go @@ -107,19 +107,9 @@ func (e *PolicyExecutor) Execute(ctx context.Context, node models.NodeDef, input return (&stubExecutor{stage: "policy", logName: "stub_policy", sleep: 200 * time.Millisecond}).Execute(ctx, node, inputs, ec) } -// DeployExecutor stubs the Deploy stage. -type DeployExecutor struct{} +// DeployExecutor lives in deploy.go (real impl) — was previously a stub. -func (e *DeployExecutor) Execute(ctx context.Context, node models.NodeDef, inputs [][]models.Item, ec *engine.ExecutionContext) (map[int][]models.Item, error) { - return (&stubExecutor{stage: "deploy", logName: "stub_deploy", sleep: 1 * time.Second}).Execute(ctx, node, inputs, ec) -} - -// PromoteExecutor stubs the Promote stage. -type PromoteExecutor struct{} - -func (e *PromoteExecutor) Execute(ctx context.Context, node models.NodeDef, inputs [][]models.Item, ec *engine.ExecutionContext) (map[int][]models.Item, error) { - return (&stubExecutor{stage: "promote", logName: "stub_promote", sleep: 500 * time.Millisecond}).Execute(ctx, node, inputs, ec) -} +// PromoteExecutor lives in promote.go (real impl) — was previously a stub. // RollbackExecutor stubs the Rollback stage. type RollbackExecutor struct{} diff --git a/pkg/executors/trigger.go b/pkg/executors/trigger.go index 7353565..f123905 100644 --- a/pkg/executors/trigger.go +++ b/pkg/executors/trigger.go @@ -9,9 +9,15 @@ import ( // TriggerExecutor is a universal trigger node that passes through the trigger data // injected by the runner. All n8n trigger types are mapped to this single executor. +// +// In addition, it stamps `fromBranch` and `toBranch` from the node's parameters +// onto every output item (without overwriting incoming values from a webhook +// payload). Downstream nodes — Build (for clone) and Promote (for the PR) — +// read those keys, so the Trigger node is the single place to change the +// branch pair for a pipeline. type TriggerExecutor struct{} -func (e *TriggerExecutor) Execute(_ context.Context, _ models.NodeDef, inputs [][]models.Item, _ *engine.ExecutionContext) (map[int][]models.Item, error) { +func (e *TriggerExecutor) Execute(_ context.Context, node models.NodeDef, inputs [][]models.Item, _ *engine.ExecutionContext) (map[int][]models.Item, error) { var items []models.Item for _, input := range inputs { items = append(items, input...) @@ -19,5 +25,25 @@ func (e *TriggerExecutor) Execute(_ context.Context, _ models.NodeDef, inputs [] if len(items) == 0 { items = []models.Item{{}} } + + fromBranch := strParam(node.Parameters, "fromBranch", "") + toBranch := strParam(node.Parameters, "toBranch", "") + + for i, it := range items { + if it == nil { + it = models.Item{} + } + if fromBranch != "" { + if _, ok := it["fromBranch"]; !ok { + it["fromBranch"] = fromBranch + } + } + if toBranch != "" { + if _, ok := it["toBranch"]; !ok { + it["toBranch"] = toBranch + } + } + items[i] = it + } return map[int][]models.Item{0: items}, nil } diff --git a/pkg/github/client.go b/pkg/github/client.go index f12fccd..426b94a 100644 --- a/pkg/github/client.go +++ b/pkg/github/client.go @@ -203,6 +203,228 @@ func (c *Client) UninstallWebhook(ctx context.Context, repo Repo, hookID int64) resp.StatusCode, truncate(string(body), 240)) } +// PullRequest is the slice of GitHub's PR shape we return upward. +type PullRequest struct { + Number int `json:"number"` + HTMLURL string `json:"html_url"` + State string `json:"state"` + Merged bool `json:"merged"` + // Head/base sha + ref are useful for audit; pulled from nested fields. + HeadSHA string `json:"-"` + HeadRef string `json:"-"` + BaseSHA string `json:"-"` + BaseRef string `json:"-"` + MergeSHA string `json:"-"` +} + +type prRaw struct { + Number int `json:"number"` + HTMLURL string `json:"html_url"` + State string `json:"state"` + Merged bool `json:"merged"` + MergedAt string `json:"merged_at"` + MergeSHA string `json:"merge_commit_sha"` + Head struct { + Ref string `json:"ref"` + SHA string `json:"sha"` + } `json:"head"` + Base struct { + Ref string `json:"ref"` + SHA string `json:"sha"` + } `json:"base"` +} + +func (r prRaw) into() PullRequest { + return PullRequest{ + Number: r.Number, HTMLURL: r.HTMLURL, + State: r.State, Merged: r.Merged, + HeadSHA: r.Head.SHA, HeadRef: r.Head.Ref, + BaseSHA: r.Base.SHA, BaseRef: r.Base.Ref, + MergeSHA: r.MergeSHA, + } +} + +// OpenPullRequest opens a PR from `head` into `base`. Returns the existing +// PR if one already exists for the same head/base pair (GitHub returns 422 +// in that case; we re-fetch via the list API). The PAT must have `repo` +// (write) permissions. +func (c *Client) OpenPullRequest(ctx context.Context, repo Repo, head, base, title, body string) (*PullRequest, error) { + payload := map[string]any{ + "title": firstNonEmpty(title, fmt.Sprintf("Promote %s → %s", head, base)), + "body": body, + "head": head, + "base": base, + } + buf, _ := json.Marshal(payload) + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, + apiBase+"/repos/"+repo.String()+"/pulls", bytes.NewReader(buf)) + if err != nil { + return nil, err + } + c.applyAuth(req) + req.Header.Set("Content-Type", "application/json") + + resp, err := c.http.Do(req) + if err != nil { + return nil, fmt.Errorf("github request failed: %w", err) + } + defer resp.Body.Close() + respBody, _ := io.ReadAll(resp.Body) + + switch resp.StatusCode { + case http.StatusCreated: + var pr prRaw + if err := json.Unmarshal(respBody, &pr); err != nil { + return nil, fmt.Errorf("decode PR response: %w", err) + } + out := pr.into() + return &out, nil + case http.StatusUnprocessableEntity: + // Most common cause: a PR already exists for this head/base. Fetch + // it so promote stays idempotent. + if existing, err := c.findPR(ctx, repo, head, base); err == nil && existing != nil { + return existing, nil + } + // Otherwise surface the original 422 (could be: no commits between + // branches, head doesn't exist, base doesn't exist, etc.). + return nil, fmt.Errorf("open PR: github 422: %s", truncate(string(respBody), 240)) + default: + return nil, fmt.Errorf("open PR: github returned %d: %s", + resp.StatusCode, truncate(string(respBody), 240)) + } +} + +// findPR queries the open PRs for a head→base match. Used to recover an +// already-open PR when OpenPullRequest's 422 means "duplicate". +func (c *Client) findPR(ctx context.Context, repo Repo, head, base string) (*PullRequest, error) { + // GitHub expects head as `:` for cross-repo search; for + // same-repo PRs the bare branch is enough. + q := url.Values{} + q.Set("state", "open") + q.Set("head", repo.Owner+":"+head) + q.Set("base", base) + url := apiBase + "/repos/" + repo.String() + "/pulls?" + q.Encode() + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, url, nil) + if err != nil { + return nil, err + } + c.applyAuth(req) + resp, err := c.http.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + body, _ := io.ReadAll(resp.Body) + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("list pulls: github %d: %s", resp.StatusCode, truncate(string(body), 240)) + } + var arr []prRaw + if err := json.Unmarshal(body, &arr); err != nil { + return nil, err + } + if len(arr) == 0 { + return nil, nil + } + out := arr[0].into() + return &out, nil +} + +// MergePullRequest merges PR #number using the API's PUT /pulls/{n}/merge. +// `commitMessage` is optional. Method may be "merge" | "squash" | "rebase". +// Returns the merge commit SHA on success. +func (c *Client) MergePullRequest(ctx context.Context, repo Repo, number int, commitMessage, method string) (string, error) { + if method == "" { + method = "merge" + } + payload := map[string]any{"merge_method": method} + if commitMessage != "" { + payload["commit_message"] = commitMessage + } + buf, _ := json.Marshal(payload) + + url := fmt.Sprintf("%s/repos/%s/pulls/%d/merge", apiBase, repo.String(), number) + req, err := http.NewRequestWithContext(ctx, http.MethodPut, url, bytes.NewReader(buf)) + if err != nil { + return "", err + } + c.applyAuth(req) + req.Header.Set("Content-Type", "application/json") + resp, err := c.http.Do(req) + if err != nil { + return "", fmt.Errorf("github request failed: %w", err) + } + defer resp.Body.Close() + body, _ := io.ReadAll(resp.Body) + if resp.StatusCode != http.StatusOK { + return "", fmt.Errorf("merge PR #%d: github %d: %s", + number, resp.StatusCode, truncate(string(body), 240)) + } + var r struct { + SHA string `json:"sha"` + Merged bool `json:"merged"` + Message string `json:"message"` + } + if err := json.Unmarshal(body, &r); err != nil { + return "", err + } + if !r.Merged { + return "", fmt.Errorf("merge not applied: %s", r.Message) + } + return r.SHA, nil +} + +// MergeBranches uses POST /repos/{owner}/{name}/merges to fast-forward +// `base` to include `head`. No PR involved. Returns the merge commit SHA, +// or "" with no error when the branches are already up to date (204). +func (c *Client) MergeBranches(ctx context.Context, repo Repo, base, head, commitMessage string) (string, error) { + payload := map[string]any{ + "base": base, + "head": head, + "commit_message": commitMessage, + } + buf, _ := json.Marshal(payload) + url := apiBase + "/repos/" + repo.String() + "/merges" + req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(buf)) + if err != nil { + return "", err + } + c.applyAuth(req) + req.Header.Set("Content-Type", "application/json") + resp, err := c.http.Do(req) + if err != nil { + return "", fmt.Errorf("github request failed: %w", err) + } + defer resp.Body.Close() + body, _ := io.ReadAll(resp.Body) + switch resp.StatusCode { + case http.StatusCreated: + var r struct{ SHA string `json:"sha"` } + if err := json.Unmarshal(body, &r); err != nil { + return "", err + } + return r.SHA, nil + case http.StatusNoContent: + // 204 = already up to date. Caller treats this as a soft success. + return "", nil + case http.StatusConflict: + return "", fmt.Errorf("merge conflict between %s and %s — resolve manually or use mode=open-pr", base, head) + default: + return "", fmt.Errorf("merge branches: github %d: %s", + resp.StatusCode, truncate(string(body), 240)) + } +} + +func firstNonEmpty(ss ...string) string { + for _, s := range ss { + if strings.TrimSpace(s) != "" { + return s + } + } + return "" +} + func (c *Client) applyAuth(req *http.Request) { if c.pat != "" { req.Header.Set("Authorization", "Bearer "+c.pat) diff --git a/pkg/orchestrator/walk.go b/pkg/orchestrator/walk.go index 7e416ca..bc4f99f 100644 --- a/pkg/orchestrator/walk.go +++ b/pkg/orchestrator/walk.go @@ -104,10 +104,22 @@ func walkDurable( } started := time.Now() - outputs, runErr := runNode(ctx, "node:"+nm, func(c context.Context) (map[int][]models.Item, error) { - c = engine.WithNodeLogger(c, nodeLogger) - return ex(c, nd, in, execCtx) - }) + // Approval is special: its executor calls Restate context methods + // (restate.Set / Clear / Awakeable) directly. Wrapping it in + // restate.Run would make those calls happen inside a Run block, + // which the SDK rejects with "Concurrent context use detected". + // So Approval runs against the parent workflow context, no Run. + var outputs map[int][]models.Item + var runErr error + if nd.Type == "flow-nodes-base.waitForApproval" { + runCtx := engine.WithNodeLogger(ctx, nodeLogger) + outputs, runErr = ex(runCtx, nd, in, execCtx) + } else { + outputs, runErr = runNode(ctx, "node:"+nm, func(c context.Context) (map[int][]models.Item, error) { + c = engine.WithNodeLogger(c, nodeLogger) + return ex(c, nd, in, execCtx) + }) + } dur := time.Since(started).Milliseconds() if runErr != nil { diff --git a/pkg/secrets/secrets.go b/pkg/secrets/secrets.go new file mode 100644 index 0000000..2e4a4b9 --- /dev/null +++ b/pkg/secrets/secrets.go @@ -0,0 +1,156 @@ +// Package secrets is the host-side encryption helper used by anything +// in Flow that persists user-supplied credentials (cloud creds, kv +// secrets, agent PATs). Master key comes from the FLOW_SECRET_KEY env +// var. Format: base64-url(nonce || ciphertext || tag), AES-256-GCM. +// +// Design choices: +// - Fail loud if the key is missing. We never silently store plaintext +// in fields the rest of the system thinks are encrypted — that's +// worse than failing to start. +// - Master key is hashed with SHA-256 so any reasonable input length +// works (32-byte hex, base64, or a passphrase). +// - Each Seal generates a fresh random 12-byte nonce. Ciphertext is +// authenticated by GCM's tag. +// +// We do NOT version the ciphertext envelope. If we ever need key +// rotation, prepend a 1-byte version + key-id and decode by version. +package secrets + +import ( + "crypto/aes" + "crypto/cipher" + "crypto/rand" + "crypto/sha256" + "encoding/base64" + "errors" + "fmt" + "io" + "os" + "sync" +) + +const envKey = "FLOW_SECRET_KEY" + +var ( + keyOnce sync.Once + keyVal []byte + keyErr error +) + +// loadKey reads FLOW_SECRET_KEY once and caches the SHA-256 of its bytes. +// Returns an error (not panic) so callers can surface it cleanly. +func loadKey() ([]byte, error) { + keyOnce.Do(func() { + raw := os.Getenv(envKey) + if raw == "" { + keyErr = fmt.Errorf("secrets: %s is not set; refusing to encrypt/decrypt", envKey) + return + } + sum := sha256.Sum256([]byte(raw)) + keyVal = sum[:] + }) + return keyVal, keyErr +} + +// SealString encrypts plaintext with AES-256-GCM and returns a +// base64-url-encoded envelope. Empty input returns empty output (so +// "no secret set" round-trips through encryption cleanly). +func SealString(plaintext string) (string, error) { + if plaintext == "" { + return "", nil + } + key, err := loadKey() + if err != nil { + return "", err + } + block, err := aes.NewCipher(key) + if err != nil { + return "", fmt.Errorf("aes new: %w", err) + } + gcm, err := cipher.NewGCM(block) + if err != nil { + return "", fmt.Errorf("gcm new: %w", err) + } + nonce := make([]byte, gcm.NonceSize()) + if _, err := io.ReadFull(rand.Reader, nonce); err != nil { + return "", fmt.Errorf("nonce: %w", err) + } + ct := gcm.Seal(nil, nonce, []byte(plaintext), nil) + envelope := append(nonce, ct...) + return base64.RawURLEncoding.EncodeToString(envelope), nil +} + +// OpenString reverses SealString. Empty input returns empty output. +// Any decode/auth failure is reported — never silently fall back to +// returning the input as plaintext. +func OpenString(envelope string) (string, error) { + if envelope == "" { + return "", nil + } + key, err := loadKey() + if err != nil { + return "", err + } + raw, err := base64.RawURLEncoding.DecodeString(envelope) + if err != nil { + return "", fmt.Errorf("base64: %w", err) + } + block, err := aes.NewCipher(key) + if err != nil { + return "", fmt.Errorf("aes new: %w", err) + } + gcm, err := cipher.NewGCM(block) + if err != nil { + return "", fmt.Errorf("gcm new: %w", err) + } + if len(raw) < gcm.NonceSize() { + return "", errors.New("ciphertext too short") + } + nonce, ct := raw[:gcm.NonceSize()], raw[gcm.NonceSize():] + pt, err := gcm.Open(nil, nonce, ct, nil) + if err != nil { + return "", fmt.Errorf("gcm open (key mismatch or tampered): %w", err) + } + return string(pt), nil +} + +// SealMap encrypts every value of m in place, returning a new map. +// Useful for kv-credential maps without writing a loop at every caller. +func SealMap(m map[string]string) (map[string]string, error) { + if m == nil { + return nil, nil + } + out := make(map[string]string, len(m)) + for k, v := range m { + sealed, err := SealString(v) + if err != nil { + return nil, fmt.Errorf("seal %q: %w", k, err) + } + out[k] = sealed + } + return out, nil +} + +// OpenMap reverses SealMap. +func OpenMap(m map[string]string) (map[string]string, error) { + if m == nil { + return nil, nil + } + out := make(map[string]string, len(m)) + for k, v := range m { + opened, err := OpenString(v) + if err != nil { + return nil, fmt.Errorf("open %q: %w", k, err) + } + out[k] = opened + } + return out, nil +} + +// IsConfigured reports whether the master key is set. Useful at startup +// so the API can refuse to mount the credentials endpoints (or warn +// loudly) when secrets won't work. +func IsConfigured() bool { + _, err := loadKey() + return err == nil +} diff --git a/pkg/storage/mongo.go b/pkg/storage/mongo.go index 72d8b6a..30946af 100644 --- a/pkg/storage/mongo.go +++ b/pkg/storage/mongo.go @@ -65,6 +65,13 @@ func (m *Mongo) Runs() RunStore { return &mongoRuns{coll: m.db.Collection("runs" // Agents returns the AgentStore backed by this Mongo connection. func (m *Mongo) Agents() AgentStore { return &mongoAgents{coll: m.db.Collection("agents")} } +// Credentials returns the global CredentialStore backed by this Mongo +// connection. Per-agent overrides live on agent.Credentials and are not +// persisted here; this collection is the org-wide pool. +func (m *Mongo) Credentials() CredentialStore { + return &mongoCredentials{coll: m.db.Collection("credentials")} +} + func (m *Mongo) ensureIndexes(ctx context.Context) error { if _, err := m.db.Collection("pipelines").Indexes().CreateMany(ctx, []mongo.IndexModel{ {Keys: bson.D{{Key: "updated_at", Value: -1}}}, @@ -82,6 +89,12 @@ func (m *Mongo) ensureIndexes(ctx context.Context) error { }); err != nil { return fmt.Errorf("agents indexes: %w", err) } + if _, err := m.db.Collection("credentials").Indexes().CreateMany(ctx, []mongo.IndexModel{ + {Keys: bson.D{{Key: "name", Value: 1}}, Options: options.Index().SetUnique(true)}, + {Keys: bson.D{{Key: "updated_at", Value: -1}}}, + }); err != nil { + return fmt.Errorf("credentials indexes: %w", err) + } return nil } @@ -447,3 +460,65 @@ func (s *mongoAgents) List(ctx context.Context) ([]*Agent, error) { } return out, nil } + +// --- credentials (global pool) ------------------------------------------- + +type mongoCredentials struct{ coll *mongo.Collection } + +func (s *mongoCredentials) Create(ctx context.Context, c *Credential) error { + if _, err := s.coll.InsertOne(ctx, c); err != nil { + // Surface the duplicate-name index violation as a typed error so + // the API layer can return 409 instead of a generic 500. + if mongo.IsDuplicateKeyError(err) { + return ErrAlreadyExists + } + return err + } + return nil +} + +func (s *mongoCredentials) GetByName(ctx context.Context, name string) (*Credential, error) { + var c Credential + if err := s.coll.FindOne(ctx, bson.M{"name": name}).Decode(&c); err != nil { + if errors.Is(err, mongo.ErrNoDocuments) { + return nil, ErrNotFound + } + return nil, err + } + return &c, nil +} + +func (s *mongoCredentials) Update(ctx context.Context, c *Credential) error { + res, err := s.coll.ReplaceOne(ctx, bson.M{"name": c.Name}, c) + if err != nil { + return err + } + if res.MatchedCount == 0 { + return ErrNotFound + } + return nil +} + +func (s *mongoCredentials) Delete(ctx context.Context, name string) error { + res, err := s.coll.DeleteOne(ctx, bson.M{"name": name}) + if err != nil { + return err + } + if res.DeletedCount == 0 { + return ErrNotFound + } + return nil +} + +func (s *mongoCredentials) List(ctx context.Context) ([]*Credential, error) { + cur, err := s.coll.Find(ctx, bson.M{}, options.Find().SetSort(bson.D{{Key: "updated_at", Value: -1}})) + if err != nil { + return nil, err + } + defer cur.Close(ctx) + var out []*Credential + if err := cur.All(ctx, &out); err != nil { + return nil, err + } + return out, nil +} diff --git a/pkg/storage/storage.go b/pkg/storage/storage.go index 6bfdfa2..53ea048 100644 --- a/pkg/storage/storage.go +++ b/pkg/storage/storage.go @@ -7,6 +7,7 @@ import ( "context" "encoding/json" "errors" + "strings" "time" ) @@ -14,6 +15,11 @@ import ( // handlers should map this to HTTP 404. var ErrNotFound = errors.New("not found") +// ErrAlreadyExists is returned by Create when a uniqueness constraint +// (e.g. credential name) would be violated. API handlers should map +// this to HTTP 409. +var ErrAlreadyExists = errors.New("already exists") + // Pipeline is the persisted shape of a pipeline (formerly "flow") that the // API stores and returns to the UI. Definition is the n8n-format JSON. type Pipeline struct { @@ -69,6 +75,49 @@ const ( AuthFailed AuthStatus = "failed" ) +// CredentialType discriminates the shape of a credential record. Only +// fields belonging to the matching type should be populated; the rest +// stay zero-valued. +type CredentialType string + +const ( + CredentialAWS CredentialType = "aws" + CredentialGCP CredentialType = "gcp" + CredentialKV CredentialType = "kv" +) + +// Credential is a named credential record attached to an agent. The +// Deploy node (and any other node that needs cloud creds) looks one up +// by Name. Secret fields are encrypted at rest with pkg/secrets; the API +// layer never returns them — only `HasSecret` flags. +// +// Fields are deliberately flat instead of `union { aws, gcp, kv }` so +// the Mongo schema stays simple and an upgrade to a new type only adds +// fields without rewriting the doc shape. +type Credential struct { + ID string `json:"id" bson:"id"` + Name string `json:"name" bson:"name"` // unique within an agent + Type CredentialType `json:"type" bson:"type"` + + // AWS — non-secret. The Flow host's own identity AssumeRoles into + // AwsCrossAccountRoleArn at deploy time. + AwsRegion string `json:"awsRegion,omitempty" bson:"aws_region,omitempty"` + AwsAccountID string `json:"awsAccountId,omitempty" bson:"aws_account_id,omitempty"` + AwsCrossAccountRoleArn string `json:"awsCrossAccountRoleArn,omitempty" bson:"aws_cross_account_role_arn,omitempty"` + + // GCP — projectId / location are non-secret. Service account JSON is. + GcpProjectID string `json:"gcpProjectId,omitempty" bson:"gcp_project_id,omitempty"` + GcpLocation string `json:"gcpLocation,omitempty" bson:"gcp_location,omitempty"` + GcpServiceAccountSealed string `json:"-" bson:"gcp_sa_sealed,omitempty"` + + // Generic key-value store. Each value is sealed independently so we + // can return a list of keys publicly without leaking values. + KvSealed map[string]string `json:"-" bson:"kv_sealed,omitempty"` + + CreatedAt time.Time `json:"createdAt" bson:"created_at"` + UpdatedAt time.Time `json:"updatedAt" bson:"updated_at"` +} + // Agent is an agent repo registered with Langship. The PAT and webhook // secret are stored server-side; the API layer scrubs them before the // record leaves the boundary (see pkg/api/agents.go). @@ -84,8 +133,26 @@ type Agent struct { AuthStatus AuthStatus `json:"authStatus,omitempty" bson:"auth_status,omitempty"` AuthCheckedAt *time.Time `json:"authCheckedAt,omitempty" bson:"auth_checked_at,omitempty"` AttachedPipelines []string `json:"attachedPipelines,omitempty" bson:"attached_pipelines,omitempty"` - CreatedAt time.Time `json:"createdAt" bson:"created_at"` - UpdatedAt time.Time `json:"updatedAt" bson:"updated_at"` + + // Named credentials — referenced by name from Deploy / future nodes. + Credentials []Credential `json:"credentials,omitempty" bson:"credentials,omitempty"` + + CreatedAt time.Time `json:"createdAt" bson:"created_at"` + UpdatedAt time.Time `json:"updatedAt" bson:"updated_at"` +} + +// LookupCredential returns the agent's credential matching name (case- +// insensitive) and an ok flag. Convenience for executors. +func (a *Agent) LookupCredential(name string) (Credential, bool) { + if a == nil { + return Credential{}, false + } + for _, c := range a.Credentials { + if strings.EqualFold(c.Name, name) { + return c, true + } + } + return Credential{}, false } // AgentStore persists agent registrations. Update mutates the entire @@ -97,3 +164,16 @@ type AgentStore interface { Delete(ctx context.Context, id string) error List(ctx context.Context) ([]*Agent, error) } + +// CredentialStore persists global (org-wide) credentials. Agents can +// override these by name with a record on agent.Credentials, but the +// global pool is the canonical place to define a credential once and +// reuse it across many agents/pipelines. Lookup is by name (the user- +// facing identifier — Deploy nodes reference creds by name, not ID). +type CredentialStore interface { + Create(ctx context.Context, c *Credential) error + GetByName(ctx context.Context, name string) (*Credential, error) + Update(ctx context.Context, c *Credential) error + Delete(ctx context.Context, name string) error + List(ctx context.Context) ([]*Credential, error) +} diff --git a/web/app/agents/new/page.tsx b/web/app/agents/new/page.tsx index b27ff38..cd5cd64 100644 --- a/web/app/agents/new/page.tsx +++ b/web/app/agents/new/page.tsx @@ -1,36 +1,102 @@ "use client"; -import { useState } from "react"; +import { useEffect, useRef, useState } from "react"; import { useRouter } from "next/navigation"; -import Link from "next/link"; -import { ArrowLeft, Eye, EyeOff, KeyRound, Save } from "lucide-react"; -import { Button } from "@/components/ui/button"; import { - Card, - CardContent, - CardDescription, - CardHeader, - CardTitle, -} from "@/components/ui/card"; + ArrowRight, + Eye, + EyeOff, + Github, + Gitlab, + KeyRound, + Sparkles, +} from "lucide-react"; + +import { Button } from "@/components/ui/button"; import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { api } from "@/lib/api"; +// ─── Templates (no-op for now) ────────────────────────────────────────────── +// Display-only cards; clicking pre-fills the URL bar with a known starter +// repo. Wiring real template cloning is a later story. +type Template = { + badge: string; + title: string; + description: string; + // repoUrl is the URL we pre-fill on click. `null` means the template + // isn't wired yet — the card renders disabled with a SOON pill. + repoUrl: string | null; +}; + +const TEMPLATES: Template[] = [ + { + badge: "LANGGRAPH", + title: "LangGraph quickstart", + description: "Stateful agent graph with tool calls + memory.", + repoUrl: "https://github.com/patel-lyzr/langraph-agent", + }, + { + badge: "CREWAI", + title: "CrewAI starter", + description: "Multi-agent crew with role-based collaboration.", + repoUrl: null, + }, + { + badge: "LANGCHAIN", + title: "LangChain agent", + description: "Classic ReAct agent with retrievers and tools.", + repoUrl: null, + }, +]; + export default function NewAgentPage() { const router = useRouter(); const [repoUrl, setRepoUrl] = useState(""); const [pat, setPat] = useState(""); const [showPat, setShowPat] = useState(false); + const [showPatField, setShowPatField] = useState(false); const [error, setError] = useState(null); const [saving, setSaving] = useState(false); + const urlInputRef = useRef(null); + const patInputRef = useRef(null); + + // Focus the URL bar after either provider/template card click so the user + // immediately has a clear next action. + useEffect(() => { + if (showPatField) { + // Give the input a beat to render before focusing. + requestAnimationFrame(() => patInputRef.current?.focus()); + } + }, [showPatField]); + + function pickTemplate(t: Template) { + if (!t.repoUrl) return; + setRepoUrl(t.repoUrl); + setShowPatField(true); + requestAnimationFrame(() => urlInputRef.current?.focus()); + } + + function continueWithGitHub() { + setShowPatField(true); + if (!repoUrl.trim()) { + requestAnimationFrame(() => urlInputRef.current?.focus()); + } else { + requestAnimationFrame(() => patInputRef.current?.focus()); + } + } + async function onSave() { setSaving(true); setError(null); try { const url = repoUrl.trim(); if (!url) throw new Error("Repository URL is required"); - await api.createAgent({ repoUrl: url, pat: pat.trim() || undefined }); + await api.createAgent({ + repoUrl: url, + pat: pat.trim() || undefined, + }); router.push("/agents"); } catch (e) { setError(e instanceof Error ? e.message : "save failed"); @@ -40,95 +106,247 @@ export default function NewAgentPage() { } return ( -
- - +
-

Add agent

+

+ Register an agent +

- Link a git repository that contains your agent code. The PAT is - stored encrypted server-side and used only for git operations. + One agent = one repo + one PAT + the pipelines that ship it.

- - - Repository - - HTTPS or SSH URL — public repos can leave the PAT blank. - - - -
- - setRepoUrl(e.target.value)} - placeholder="https://github.com/org/agent-repo.git" - spellCheck={false} - autoComplete="off" - /> -

- Name is auto-derived from the repo path (e.g.{" "} - org/repo). -

-
+ {/* URL bar -------------------------------------------------------- */} +
+ + setRepoUrl(e.target.value)} + placeholder="Paste a GitHub repo URL — or pick a template below" + spellCheck={false} + autoComplete="off" + className="h-10 border-0 bg-transparent text-sm focus-visible:ring-0" + /> + +
+ {/* Two-column: provider + templates ------------------------------- */} +
+ {/* Import Git Repository ---------------------------------------- */} +
+
+ Import Git Repository +
+

+ Pick a provider. We use a PAT to install a webhook on your repo so + push/PR events trigger pipelines. +

- -
+ + + +
+

+ Only GitHub is wired up in v0. GitLab and Bitbucket are next. +

+
+ + {/* Clone Template ----------------------------------------------- */} +
+
+
Clone Template
+ + Framework starters + +
+
+ {TEMPLATES.map((t) => { + const disabled = !t.repoUrl; + return ( + + ); + })} +
+
+ More framework starters coming soon. Want one added? Open an issue + on the Langship repo. +
+
+
+ + {/* PAT entry + Save (revealed after a provider/template is chosen) - */} + {showPatField && ( +
+
Connect repo
+

+ HTTPS URL above; PAT below. Public repos can leave the PAT blank. +

+
+
+ setPat(e.target.value)} - placeholder="ghp_… or glpat_…" + id="repoUrl" + value={repoUrl} + onChange={(e) => setRepoUrl(e.target.value)} + placeholder="https://github.com/org/agent-repo.git" spellCheck={false} - autoComplete="new-password" - className="pr-10 font-mono text-xs" + autoComplete="off" /> - +

+ Agent name is auto-derived from the repo path (e.g.{" "} + org/repo). +

+
+ +
+ +
+ setPat(e.target.value)} + placeholder="ghp_… or glpat_…" + spellCheck={false} + autoComplete="new-password" + className="pr-10 font-mono text-xs" + /> + +
+

+ Required for private repos. Scope:{" "} + repo read access is enough. + Stored encrypted server-side; never returned by the API after + save. +

+
+ + {error && ( +

+ {error} +

+ )} + +
+ +

- Required for private repos. Scope:{" "} - repo read access is enough. - Never returned by the API after save. + After saving, open the agent to attach pipelines and (if you’ll + use the Deploy node) override credentials.

+
+ )} - {error && ( -

- {error} -

- )} - -
- - -
- - + {/* Empty agent --------------------------------------------------- */} +
+
+
Empty agent
+

+ Skip git for now and configure the connection later. Coming soon + — for now an agent must be registered with a repo + PAT. +

+
+ +
); } + +function ProviderDisabled({ + icon: Icon, + label, +}: { + icon: React.ComponentType<{ className?: string }>; + label: string; +}) { + return ( +
+ + {label} + + SOON + +
+ ); +} + +// Bitbucket isn't in lucide-react. Tiny inline SVG to match the visual. +function BitbucketIcon({ className }: { className?: string }) { + return ( + + ); +} diff --git a/web/app/agents/view/page.tsx b/web/app/agents/view/page.tsx index 1a4c82d..7a67ca0 100644 --- a/web/app/agents/view/page.tsx +++ b/web/app/agents/view/page.tsx @@ -9,6 +9,7 @@ import { ExternalLink, Github, KeyRound, + Lock, Play, Plus, Trash2, @@ -25,11 +26,16 @@ import { CardTitle, } from "@/components/ui/card"; import { Badge } from "@/components/ui/badge"; +import { + CredentialForm, + CredentialRow, +} from "@/components/credentials/credential-form"; import { api, type Agent, type AuthStatus, type FlowSummary, + type PublicCredential, type Run, type ServerConfig, } from "@/lib/api"; @@ -512,6 +518,17 @@ function AgentDetail() {
+ {/* Credentials ------------------------------------------------------ */} + { + // refetch agent so the credentials list updates + const a = await api.getAgent(id); + setAgent(a); + }} + /> + {/* Recent runs ------------------------------------------------------ */} @@ -629,3 +646,133 @@ function RunStatus({ status }: { status: string }) { // Avoid unused-import lint when the symbol is referenced only by type. void KeyRound; void ExternalLink; + + +// ─── Credentials ──────────────────────────────────────────────────────────── + +type CredentialsSectionProps = { + agentId: string; + credentials: PublicCredential[]; + onChanged: () => void | Promise; +}; + +function CredentialsSection({ agentId, credentials, onChanged }: CredentialsSectionProps) { + const [adding, setAdding] = useState(false); + const [editingName, setEditingName] = useState(null); + const [error, setError] = useState(null); + const [globals, setGlobals] = useState([]); + + // Pull globals once so we can show inherited rows alongside the + // per-agent overrides. Refreshed when `onChanged` re-fetches the agent + // (cheap — credentials list is small). + useEffect(() => { + let cancelled = false; + api.listGlobalCredentials().then((g) => { + if (!cancelled) setGlobals(g); + }).catch(() => {}); + return () => { cancelled = true; }; + }, [credentials]); + + // Globals shadowed by an agent override: hide them from the inherited + // list; the override row is the source of truth. + const overrideNames = new Set(credentials.map((c) => c.name.toLowerCase())); + const inherited = globals.filter((g) => !overrideNames.has(g.name.toLowerCase())); + + async function handleDelete(name: string) { + if (!confirm(`Delete agent override "${name}"? The pipeline will fall back to the global credential of the same name (if any).`)) return; + try { + await api.deleteCredential(agentId, name); + await onChanged(); + } catch (e) { + setError(e instanceof Error ? e.message : "delete failed"); + } + } + + return ( + + +
+ + + Credentials + + + Cloud creds available to nodes for this agent. Globals defined on + the Credentials page are + inherited; add an override here to specialize a credential for + this agent only. + +
+ +
+ + {error && ( +

+ {error} +

+ )} + + {!adding && credentials.length === 0 && inherited.length === 0 && ( +

+ No credentials available. Add a global one on the{" "} + Credentials page{" "} + or an agent-specific override here. +

+ )} + + {credentials.map((c) => ( +
+ {editingName === c.name ? ( + setEditingName(null)} + onSubmit={async (body) => { + await api.updateCredential(agentId, c.name, body); + setEditingName(null); + await onChanged(); + }} + onError={setError} + /> + ) : ( + setEditingName(c.name)} + onDelete={() => handleDelete(c.name)} + /> + )} +
+ ))} + + {inherited.map((c) => ( +
+ { /* edit globals on the global page */ }} + onDelete={() => { /* deletes go through global page */ }} + /> +
+ ))} + + {adding && ( +
+ setAdding(false)} + onSubmit={async (body) => { + await api.createCredential(agentId, body); + setAdding(false); + await onChanged(); + }} + onError={setError} + /> +
+ )} +
+
+ ); +} diff --git a/web/app/approvals/page.tsx b/web/app/approvals/page.tsx index 120c6b3..7062805 100644 --- a/web/app/approvals/page.tsx +++ b/web/app/approvals/page.tsx @@ -1,10 +1,240 @@ -import { EmptySection } from "@/components/empty-section"; +"use client"; + +import { useCallback, useEffect, useRef, useState } from "react"; +import Link from "next/link"; +import { + ArrowRight, + Check, + Pause, + RefreshCw, + X, +} from "lucide-react"; +import { Button } from "@/components/ui/button"; +import { + Card, + CardContent, + CardDescription, + CardHeader, + CardTitle, +} from "@/components/ui/card"; +import { Badge } from "@/components/ui/badge"; +import { Textarea } from "@/components/ui/textarea"; +import { api, type ExecutionStatus, type Run } from "@/lib/api"; +import { formatDate } from "@/lib/utils"; + +interface PendingApproval { + awakeable_id: string; + node?: string; + context?: { reason?: string; node?: string; [k: string]: unknown }; +} + +interface PendingItem { + run: Run; + approval: PendingApproval; +} + +// /approvals lists every run that is currently parked on a HITL Approval +// node. The orchestrator stores `pending_approval` as Restate KV and +// surfaces it via GetExecution; we fan out per "running" run, keep only +// the ones that have a pending_approval set, and let the user resolve +// them in one click. +export default function ApprovalsInboxPage() { + const [items, setItems] = useState(null); + const [error, setError] = useState(null); + const [busyId, setBusyId] = useState(null); + + const loadRef = useRef(0); + + const load = useCallback(async () => { + const tag = ++loadRef.current; + try { + const runs = await api.listRuns({ limit: 50 }); + const candidates = runs.filter((r) => /running|pending|paused|waiting/i.test(r.status)); + const checked = await Promise.all( + candidates.map(async (r) => { + try { + const s = (await api.getExecution(r.id)) as ExecutionStatus & { + pending_approval?: PendingApproval; + }; + const pa = s.pending_approval; + if (pa && pa.awakeable_id) { + return { run: r, approval: pa } as PendingItem; + } + } catch { + /* ignore — run may have advanced */ + } + return null; + }) + ); + // Only commit if a newer load hasn't started. + if (tag !== loadRef.current) return; + setItems(checked.filter(Boolean) as PendingItem[]); + setError(null); + } catch (e) { + if (tag !== loadRef.current) return; + setError(e instanceof Error ? e.message : "load failed"); + } + }, []); + + useEffect(() => { + load(); + // 6s — same reasoning as executions/view: avoid racing the workflow + // SDK's shared-handler with parallel reads. + const t = setInterval(load, 6000); + // Also reload when ANY new run is created (covers the case where a + // freshly-triggered run instantly parks on an Approval node). + const es = new EventSource(api.runsStreamURL()); + es.onmessage = () => load(); + es.onerror = () => {}; + return () => { + clearInterval(t); + es.close(); + }; + }, [load]); + + async function decide(it: PendingItem, approved: boolean, reason: string) { + setBusyId(it.run.id); + try { + await api.resumeExecution(it.run.id, { + awakeable_id: it.approval.awakeable_id, + data: { approved, reason }, + }); + await load(); + } catch (e) { + setError(e instanceof Error ? e.message : "resume failed"); + } finally { + setBusyId(null); + } + } -export default function ApprovalsPage() { return ( - +
+
+
+

Approvals

+

+ Runs paused on a HITL Approval node, awaiting a decision. +

+
+ +
+ + {error && ( + + {error} + + )} + + {items === null ? ( +

Loading…

+ ) : items.length === 0 ? ( + + + +
No approvals waiting
+

+ When a pipeline hits a Wait-for-approval node it’ll show up here + for a one-click decision. +

+
+
+ ) : ( +
+ {items.map((it) => ( + + ))} +
+ )} +
+ ); +} + +function ApprovalCard({ + item, + busy, + onDecide, +}: { + item: PendingItem; + busy: boolean; + onDecide: (it: PendingItem, approved: boolean, reason: string) => void; +}) { + const [reason, setReason] = useState(""); + const reasonHint = + typeof item.approval.context?.reason === "string" + ? (item.approval.context.reason as string) + : undefined; + const nodeName = + item.approval.node || + (typeof item.approval.context?.node === "string" + ? (item.approval.context.node as string) + : "Approval"); + + return ( + + +
+
+ + + {item.run.pipelineName || "Pipeline"} + {nodeName} + + + {item.run.id} + +
+ +
+ {reasonHint && ( +

{reasonHint}

+ )} +
+ paused since {formatDate(item.run.startedAt)} +
+
+ +
+