diff --git a/cmd/gateway/admin_api.go b/cmd/gateway/admin_api.go index 94192de..e2b035a 100644 --- a/cmd/gateway/admin_api.go +++ b/cmd/gateway/admin_api.go @@ -31,6 +31,10 @@ func newAdminAPIHandler() http.HandlerFunc { PluginArtifacts: handleAdminPluginArtifacts, PluginArtifact: handleAdminPluginArtifact, + PluginSources: handleAdminPluginSources, + PluginBuilds: handleAdminPluginBuilds, + PluginBuild: handleAdminPluginBuild, + PluginGC: handleAdminPluginGC, PluginsList: handleAdminPluginsList, PluginItem: handleAdminPluginItem, PluginAction: handleAdminPluginAction, diff --git a/cmd/gateway/admin_plugin_handlers.go b/cmd/gateway/admin_plugin_handlers.go index 65d8826..bb41c3d 100644 --- a/cmd/gateway/admin_plugin_handlers.go +++ b/cmd/gateway/admin_plugin_handlers.go @@ -6,6 +6,7 @@ import ( "net/http" "os" "path/filepath" + "strconv" "strings" "github.com/tursom/mc-gateway/internal/adminhttp" @@ -53,6 +54,179 @@ func handleAdminPluginArtifacts(w http.ResponseWriter, r *http.Request) { } } +func handleAdminPluginSources(w http.ResponseWriter, r *http.Request) { + session, ok := requireRole(w, r, adminRoleMember) + if !ok { + return + } + if pluginsManager == nil { + adminhttp.WriteAPIError(w, http.StatusServiceUnavailable, "plugin manager is not initialized") + return + } + switch r.Method { + case http.MethodGet: + artifacts, err := pluginsManager.ListArtifacts(r.Context(), r.URL.Query().Get("plugin_id")) + if err != nil { + adminhttp.WriteAPIError(w, http.StatusInternalServerError, err.Error()) + return + } + var sources []pluginmanager.ArtifactRecord + for _, artifact := range artifacts { + if artifact.ArtifactType == pluginmanager.ArtifactTypeSource { + sources = append(sources, artifact) + } + } + adminhttp.WriteJSON(w, http.StatusOK, map[string]any{"sources": sources}) + case http.MethodPost: + if session.Role != adminRoleAdmin { + adminhttp.WriteAPIError(w, http.StatusForbidden, "forbidden") + return + } + source, err := receivePluginSource(r, session.Username) + if err != nil { + recordAudit(r.Context(), session.Username, adminhttp.RequestSourceIP(r), "plugin_source_upload", "plugin_source", "", false, err.Error()) + adminhttp.WriteAPIError(w, http.StatusBadRequest, err.Error()) + return + } + recordAuditMetadata(r.Context(), session.Username, adminhttp.RequestSourceIP(r), "plugin_source_upload", "plugin_source", source.ID, true, "source package uploaded", map[string]any{ + "plugin_id": source.PluginID, + "version": source.Version, + "source_sha256": source.SHA256, + }) + adminhttp.WriteJSON(w, http.StatusCreated, map[string]any{"source": source}) + default: + adminhttp.WriteAPIError(w, http.StatusMethodNotAllowed, "method not allowed") + } +} + +func handleAdminPluginBuilds(w http.ResponseWriter, r *http.Request) { + session, ok := requireRole(w, r, adminRoleMember) + if !ok { + return + } + if pluginsManager == nil { + adminhttp.WriteAPIError(w, http.StatusServiceUnavailable, "plugin manager is not initialized") + return + } + switch r.Method { + case http.MethodGet: + builds, err := pluginsManager.ListBuilds(r.Context(), r.URL.Query().Get("plugin_id")) + if err != nil { + adminhttp.WriteAPIError(w, http.StatusInternalServerError, err.Error()) + return + } + adminhttp.WriteJSON(w, http.StatusOK, map[string]any{"builds": builds}) + case http.MethodPost: + if session.Role != adminRoleAdmin { + adminhttp.WriteAPIError(w, http.StatusForbidden, "forbidden") + return + } + var req pluginmanager.BuildRequest + if !adminhttp.DecodeJSONRequest(w, r, &req) { + return + } + build, err := pluginsManager.CreateBuild(r.Context(), session.Username, req) + if err != nil { + recordAudit(r.Context(), session.Username, adminhttp.RequestSourceIP(r), "plugin_source_build", "plugin_build", req.SourceID, false, err.Error()) + adminhttp.WriteAPIError(w, http.StatusBadRequest, err.Error()) + return + } + recordAuditMetadata(r.Context(), session.Username, adminhttp.RequestSourceIP(r), "plugin_source_build_queue", "plugin_build", strconv.FormatInt(build.ID, 10), true, "build queued", map[string]any{ + "plugin_id": build.PluginID, + "source_sha256": build.SourceSHA256, + "artifact_sha256": build.ArtifactSHA256, + "builder_type": build.BuilderType, + "go_version": build.GoVersion, + }) + adminhttp.WriteJSON(w, http.StatusCreated, map[string]any{"build": build}) + default: + adminhttp.WriteAPIError(w, http.StatusMethodNotAllowed, "method not allowed") + } +} + +func handleAdminPluginBuild(w http.ResponseWriter, r *http.Request, rawSegment string) { + session, ok := requireRole(w, r, adminRoleMember) + if !ok { + return + } + if pluginsManager == nil { + adminhttp.WriteAPIError(w, http.StatusServiceUnavailable, "plugin manager is not initialized") + return + } + parts := strings.Split(rawSegment, "/") + id, err := strconv.ParseInt(parts[0], 10, 64) + if err != nil { + adminhttp.WriteAPIError(w, http.StatusBadRequest, "invalid build id") + return + } + if len(parts) == 1 && r.Method == http.MethodGet { + build, err := pluginsManager.Build(r.Context(), id) + if err != nil { + adminhttp.WriteAPIError(w, http.StatusNotFound, err.Error()) + return + } + adminhttp.WriteJSON(w, http.StatusOK, map[string]any{"build": build}) + return + } + if session.Role != adminRoleAdmin { + adminhttp.WriteAPIError(w, http.StatusForbidden, "forbidden") + return + } + if len(parts) != 2 || r.Method != http.MethodPost { + adminhttp.WriteAPIError(w, http.StatusMethodNotAllowed, "method not allowed") + return + } + action, err := adminhttp.PathSegment(parts[1]) + if err != nil { + adminhttp.WriteAPIError(w, http.StatusBadRequest, err.Error()) + return + } + var build pluginmanager.BuildRecord + switch action { + case "run": + build, err = pluginsManager.RunBuild(r.Context(), session.Username, id) + case "cancel": + build, err = pluginsManager.CancelBuild(r.Context(), session.Username, id) + case "retry": + build, err = pluginsManager.RetryBuild(r.Context(), session.Username, id) + default: + adminhttp.WriteAPIError(w, http.StatusBadRequest, "unknown build action") + return + } + if err != nil { + recordAudit(r.Context(), session.Username, adminhttp.RequestSourceIP(r), "plugin_build_"+action, "plugin_build", strconv.FormatInt(id, 10), false, err.Error()) + adminhttp.WriteAPIError(w, http.StatusBadRequest, err.Error()) + return + } + recordAudit(r.Context(), session.Username, adminhttp.RequestSourceIP(r), "plugin_build_"+action, "plugin_build", strconv.FormatInt(id, 10), true, "build "+action+" succeeded") + adminhttp.WriteJSON(w, http.StatusOK, map[string]any{"build": build}) +} + +func handleAdminPluginGC(w http.ResponseWriter, r *http.Request) { + session, ok := requireRole(w, r, adminRoleMember) + if !ok { + return + } + if pluginsManager == nil { + adminhttp.WriteAPIError(w, http.StatusServiceUnavailable, "plugin manager is not initialized") + return + } + dryRun := r.Method == http.MethodGet + if r.Method == http.MethodPost && session.Role != adminRoleAdmin { + adminhttp.WriteAPIError(w, http.StatusForbidden, "forbidden") + return + } + candidates, err := pluginsManager.RunGC(r.Context(), session.Username, dryRun) + if err != nil { + adminhttp.WriteAPIError(w, http.StatusInternalServerError, err.Error()) + return + } + adminhttp.WriteJSON(w, http.StatusOK, map[string]any{ + "dry_run": dryRun, + "candidates": candidates, + }) +} + func handleAdminPluginArtifact(w http.ResponseWriter, r *http.Request, rawArtifactID string) { if _, ok := requireRole(w, r, adminRoleMember); !ok { return @@ -271,7 +445,50 @@ func receivePluginArtifact(r *http.Request, actor string) (pluginmanager.Artifac return pluginmanager.ArtifactRecord{}, err } - return pluginsManager.UploadArtifact(r.Context(), pluginmanager.ArtifactUpload{ + upload := pluginmanager.ArtifactUpload{ + SourcePath: tmpPath, + FileName: filepath.Base(header.Filename), + Actor: actor, + } + manifest, err := readPackageManifest(tmpPath) + if err != nil { + return pluginmanager.ArtifactRecord{}, err + } + if manifest.ArtifactType == pluginmanager.ArtifactTypeSource { + return pluginsManager.UploadSource(r.Context(), upload) + } + return pluginsManager.UploadArtifact(r.Context(), upload) +} + +func receivePluginSource(r *http.Request, actor string) (pluginmanager.ArtifactRecord, error) { + if err := r.ParseMultipartForm(64 << 20); err != nil { + return pluginmanager.ArtifactRecord{}, err + } + file, header, err := r.FormFile("artifact") + if err != nil { + file, header, err = r.FormFile("source") + } + if err != nil { + return pluginmanager.ArtifactRecord{}, err + } + defer file.Close() + + tmp, err := os.CreateTemp("", "mc-gateway-plugin-source-*.mcgp") + if err != nil { + return pluginmanager.ArtifactRecord{}, err + } + tmpPath := tmp.Name() + defer os.Remove(tmpPath) + defer tmp.Close() + + if _, err := tmp.ReadFrom(file); err != nil { + return pluginmanager.ArtifactRecord{}, err + } + if err := tmp.Close(); err != nil { + return pluginmanager.ArtifactRecord{}, err + } + + return pluginsManager.UploadSource(r.Context(), pluginmanager.ArtifactUpload{ SourcePath: tmpPath, FileName: filepath.Base(header.Filename), Actor: actor, diff --git a/cmd/gateway/plugin_cli.go b/cmd/gateway/plugin_cli.go index 8d409f4..018afa0 100644 --- a/cmd/gateway/plugin_cli.go +++ b/cmd/gateway/plugin_cli.go @@ -2,12 +2,15 @@ package main import ( "archive/zip" + "context" + "database/sql" "encoding/json" "fmt" "io" "os" "path/filepath" + "github.com/tursom/mc-gateway/internal/admindb" "github.com/tursom/mc-gateway/internal/pluginmanager" ) @@ -16,7 +19,7 @@ func runPluginCLI(args []string) (bool, int) { return false, 0 } if len(args) < 3 { - fmt.Fprintln(os.Stderr, "usage: gateway plugin inspect|validate|compat ") + fmt.Fprintln(os.Stderr, "usage: gateway plugin inspect|validate|compat|source-validate | source-build [out.mcgp]") return true, 2 } @@ -55,12 +58,159 @@ func runPluginCLI(args []string) (bool, int) { fmt.Fprintf(os.Stdout, "ok plugin=%s version=%s sha256=%s api=%s go=%s %s/%s\n", artifact.PluginID, artifact.Version, artifact.SHA256, artifact.APIVersion, artifact.GoVersion, artifact.GOOS, artifact.GOARCH) return true, 0 + case "source-validate": + tmpRoot, err := os.MkdirTemp("", "mcgp-source-cli-*") + if err != nil { + fmt.Fprintln(os.Stderr, err) + return true, 1 + } + defer os.RemoveAll(tmpRoot) + store := pluginmanager.NewArtifactStore(tmpRoot) + source, err := store.ValidateAndStoreSource(pluginmanager.ArtifactUpload{ + SourcePath: packagePath, + FileName: filepath.Base(packagePath), + Actor: "cli", + }) + if err != nil { + fmt.Fprintln(os.Stderr, err) + return true, 1 + } + fmt.Fprintf(os.Stdout, "ok source plugin=%s version=%s source_sha256=%s api=%s go=%s %s/%s\n", + source.PluginID, source.Version, source.SHA256, source.APIVersion, source.GoVersion, source.GOOS, source.GOARCH) + return true, 0 + case "source-build": + outPath := "" + if len(args) >= 4 { + outPath = args[3] + } + build, out, err := buildSourcePackageForCLI(packagePath, outPath) + if err != nil { + fmt.Fprintln(os.Stderr, err) + return true, 1 + } + fmt.Fprintf(os.Stdout, "ok build=%d plugin=%s source_sha256=%s artifact_sha256=%s builder=%s go=%s status=%s out=%s\n", + build.ID, build.PluginID, build.SourceSHA256, build.ArtifactSHA256, build.BuilderType, build.GoVersion, build.Status, out) + return true, 0 default: fmt.Fprintf(os.Stderr, "unknown plugin command %q\n", command) return true, 2 } } +func buildSourcePackageForCLI(packagePath, outPath string) (pluginmanager.BuildRecord, string, error) { + tmpRoot, err := os.MkdirTemp("", "mcgp-source-build-cli-*") + if err != nil { + return pluginmanager.BuildRecord{}, "", err + } + defer os.RemoveAll(tmpRoot) + db, err := openPluginCLIDB(filepath.Join(tmpRoot, "plugins.db")) + if err != nil { + return pluginmanager.BuildRecord{}, "", err + } + defer db.Close() + manager := pluginmanager.New(pluginmanager.Options{ + DB: db, + ArtifactRoot: filepath.Join(tmpRoot, "artifacts"), + }) + source, err := manager.UploadSource(context.Background(), pluginmanager.ArtifactUpload{ + SourcePath: packagePath, + FileName: filepath.Base(packagePath), + Actor: "cli", + }) + if err != nil { + return pluginmanager.BuildRecord{}, "", err + } + builds, err := manager.ListBuilds(context.Background(), source.PluginID) + if err != nil { + return pluginmanager.BuildRecord{}, "", err + } + var queued pluginmanager.BuildRecord + for _, build := range builds { + if build.SourceID == source.ID && build.Status == pluginmanager.BuildStatusQueued { + queued = build + break + } + } + if queued.ID == 0 { + return pluginmanager.BuildRecord{}, "", fmt.Errorf("source upload did not create a queued build") + } + build, err := manager.RunBuild(context.Background(), "cli", queued.ID) + if err != nil { + return pluginmanager.BuildRecord{}, "", err + } + if build.Status != pluginmanager.BuildStatusSucceeded { + if build.LogSummary != "" { + return build, "", fmt.Errorf("source build status %s: %s\n%s", build.Status, build.Error, build.LogSummary) + } + return build, "", fmt.Errorf("source build status %s: %s", build.Status, build.Error) + } + artifact, err := manager.Artifact(context.Background(), build.ArtifactID) + if err != nil { + return build, "", err + } + if outPath == "" { + outPath = filepath.Join(filepath.Dir(packagePath), artifact.PluginID+"-built.mcgp") + } + if err := packageBinaryArtifact(artifact, outPath); err != nil { + return build, "", err + } + return build, outPath, nil +} + +func packageBinaryArtifact(artifact pluginmanager.ArtifactRecord, outPath string) error { + if err := os.MkdirAll(filepath.Dir(outPath), 0755); err != nil { + return err + } + out, err := os.Create(outPath) + if err != nil { + return err + } + defer out.Close() + zw := zip.NewWriter(out) + for _, entry := range []struct { + name string + path string + }{ + {name: "manifest.json", path: filepath.Join(filepath.Dir(artifact.FilePath), "manifest.json")}, + {name: pluginmanager.RuntimeEntry, path: artifact.FilePath}, + } { + if err := addZipFile(zw, entry.name, entry.path); err != nil { + zw.Close() + return err + } + } + if err := zw.Close(); err != nil { + return err + } + return out.Close() +} + +func addZipFile(zw *zip.Writer, name, path string) error { + writer, err := zw.Create(name) + if err != nil { + return err + } + file, err := os.Open(path) + if err != nil { + return err + } + defer file.Close() + _, err = io.Copy(writer, file) + return err +} + +func openPluginCLIDB(path string) (*sql.DB, error) { + db, err := admindb.Open(path) + if err != nil { + return nil, err + } + if err := admindb.Migrate(db); err != nil { + db.Close() + return nil, err + } + return db, nil +} + func readPackageManifest(packagePath string) (pluginmanager.Manifest, error) { reader, err := zip.OpenReader(packagePath) if err != nil { diff --git a/examples/plugins/mc-auth-proxy/README.md b/examples/plugins/mc-auth-proxy/README.md index 50e6454..a2cb1b5 100644 --- a/examples/plugins/mc-auth-proxy/README.md +++ b/examples/plugins/mc-auth-proxy/README.md @@ -15,7 +15,14 @@ Build and package: ./build.sh ``` -The package is written to `dist/mc-auth-proxy.mcgp`. +The binary package is written to `dist/mc-auth-proxy.mcgp`; the source package +is written to `dist/mc-auth-proxy-source.mcgp`. + +Build the source package through the gateway builder: + +```sh +go run ../../../cmd/gateway plugin source-build dist/mc-auth-proxy-source.mcgp dist/mc-auth-proxy-built.mcgp +``` Example config JSON: diff --git a/examples/plugins/mc-auth-proxy/build.sh b/examples/plugins/mc-auth-proxy/build.sh index 5d2e630..7c9168e 100755 --- a/examples/plugins/mc-auth-proxy/build.sh +++ b/examples/plugins/mc-auth-proxy/build.sh @@ -10,3 +10,14 @@ cp README.md dist/README.md rm -f mc-auth-proxy.mcgp zip -q mc-auth-proxy.mcgp manifest.json plugin.so README.md ) +rm -rf dist/source-package +mkdir -p dist/source-package/cmd/render-manifest +ARTIFACT_TYPE=source go run ./cmd/render-manifest > dist/source-package/manifest.json +cp main.go main_test.go go.mod README.md dist/source-package/ +cp cmd/render-manifest/main.go dist/source-package/cmd/render-manifest/main.go +go mod vendor -o dist/source-package/vendor +( + cd dist/source-package + rm -f ../mc-auth-proxy-source.mcgp + zip -qr ../mc-auth-proxy-source.mcgp manifest.json main.go main_test.go go.mod README.md cmd/render-manifest/main.go vendor +) diff --git a/examples/plugins/mc-auth-proxy/cmd/render-manifest/main.go b/examples/plugins/mc-auth-proxy/cmd/render-manifest/main.go index 0132ae3..8347e14 100644 --- a/examples/plugins/mc-auth-proxy/cmd/render-manifest/main.go +++ b/examples/plugins/mc-auth-proxy/cmd/render-manifest/main.go @@ -7,13 +7,17 @@ import ( ) func main() { + artifactType := os.Getenv("ARTIFACT_TYPE") + if artifactType == "" { + artifactType = "binary" + } manifest := map[string]any{ "schema_version": "mc-gateway.plugin/v1", "id": "mc-auth-proxy", "name": "Minecraft Auth Proxy", "version": "0.1.0", "description": "Protocol-proxy example that reads handshake/login start and returns a login disconnect fixture.", - "artifact_type": "binary", + "artifact_type": artifactType, "runtime": map[string]any{ "type": "go-plugin", "entry": "plugin.so", @@ -76,6 +80,17 @@ func main() { }, }, } + if artifactType == "source" { + manifest["build"] = map[string]any{ + "type": "go", + "entry": ".", + "go_version": runtime.Version(), + "cgo_enabled": true, + "tags": []string{}, + "vendor_required": false, + "output": "plugin.so", + } + } encoder := json.NewEncoder(os.Stdout) encoder.SetIndent("", " ") if err := encoder.Encode(manifest); err != nil { diff --git a/examples/plugins/upstream-rewrite/README.md b/examples/plugins/upstream-rewrite/README.md index 5cdedb7..efe1b07 100644 --- a/examples/plugins/upstream-rewrite/README.md +++ b/examples/plugins/upstream-rewrite/README.md @@ -11,7 +11,14 @@ Build and package: ./build.sh ``` -The package is written to `dist/upstream-rewrite.mcgp`. +The binary package is written to `dist/upstream-rewrite.mcgp`; the source +package is written to `dist/upstream-rewrite-source.mcgp`. + +Build the source package through the gateway builder: + +```sh +go run ../../../cmd/gateway plugin source-build dist/upstream-rewrite-source.mcgp dist/upstream-rewrite-built.mcgp +``` Example config JSON: diff --git a/examples/plugins/upstream-rewrite/build.sh b/examples/plugins/upstream-rewrite/build.sh index 674ea56..a8fed58 100755 --- a/examples/plugins/upstream-rewrite/build.sh +++ b/examples/plugins/upstream-rewrite/build.sh @@ -10,3 +10,14 @@ cp README.md dist/README.md rm -f upstream-rewrite.mcgp zip -q upstream-rewrite.mcgp manifest.json plugin.so README.md ) +rm -rf dist/source-package +mkdir -p dist/source-package/cmd/render-manifest +ARTIFACT_TYPE=source go run ./cmd/render-manifest > dist/source-package/manifest.json +cp main.go go.mod README.md dist/source-package/ +cp cmd/render-manifest/main.go dist/source-package/cmd/render-manifest/main.go +go mod vendor -o dist/source-package/vendor +( + cd dist/source-package + rm -f ../upstream-rewrite-source.mcgp + zip -qr ../upstream-rewrite-source.mcgp manifest.json main.go go.mod README.md cmd/render-manifest/main.go vendor +) diff --git a/examples/plugins/upstream-rewrite/cmd/render-manifest/main.go b/examples/plugins/upstream-rewrite/cmd/render-manifest/main.go index 30e10a2..cbe47a0 100644 --- a/examples/plugins/upstream-rewrite/cmd/render-manifest/main.go +++ b/examples/plugins/upstream-rewrite/cmd/render-manifest/main.go @@ -7,13 +7,17 @@ import ( ) func main() { + artifactType := os.Getenv("ARTIFACT_TYPE") + if artifactType == "" { + artifactType = "binary" + } manifest := map[string]any{ "schema_version": "mc-gateway.plugin/v1", "id": "upstream-rewrite", "name": "Upstream Rewrite", "version": "0.1.0", "description": "Rewrite selected upstream targets before dialing.", - "artifact_type": "binary", + "artifact_type": artifactType, "runtime": map[string]any{ "type": "go-plugin", "entry": "plugin.so", @@ -47,6 +51,17 @@ func main() { "required": []string{"upstream"}, }, } + if artifactType == "source" { + manifest["build"] = map[string]any{ + "type": "go", + "entry": ".", + "go_version": runtime.Version(), + "cgo_enabled": true, + "tags": []string{}, + "vendor_required": false, + "output": "plugin.so", + } + } encoder := json.NewEncoder(os.Stdout) encoder.SetIndent("", " ") if err := encoder.Encode(manifest); err != nil { diff --git a/internal/admindb/db.go b/internal/admindb/db.go index e0503a4..9a19969 100644 --- a/internal/admindb/db.go +++ b/internal/admindb/db.go @@ -144,6 +144,44 @@ CREATE TABLE IF NOT EXISTS plugin_operations ( created_at INTEGER NOT NULL ); +CREATE TABLE IF NOT EXISTS plugin_builds ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + plugin_id TEXT NOT NULL DEFAULT '', + source_id TEXT NOT NULL DEFAULT '', + artifact_id TEXT NOT NULL DEFAULT '', + status TEXT NOT NULL DEFAULT 'queued', + builder_type TEXT NOT NULL DEFAULT '', + builder_image TEXT NOT NULL DEFAULT '', + builder_version TEXT NOT NULL DEFAULT '', + go_version TEXT NOT NULL DEFAULT '', + go_os TEXT NOT NULL DEFAULT '', + go_arch TEXT NOT NULL DEFAULT '', + go_amd64 TEXT NOT NULL DEFAULT '', + go_arm64 TEXT NOT NULL DEFAULT '', + cgo_enabled TEXT NOT NULL DEFAULT '', + build_tags TEXT NOT NULL DEFAULT '', + sdk_module TEXT NOT NULL DEFAULT '', + sdk_version TEXT NOT NULL DEFAULT '', + go_proxy TEXT NOT NULL DEFAULT '', + go_no_sumdb TEXT NOT NULL DEFAULT '', + go_private TEXT NOT NULL DEFAULT '', + vendor_required INTEGER NOT NULL DEFAULT 0, + source_sha256 TEXT NOT NULL DEFAULT '', + artifact_sha256 TEXT NOT NULL DEFAULT '', + module_summary_json TEXT NOT NULL DEFAULT '[]', + go_version_m_json TEXT NOT NULL DEFAULT '{}', + abi_fingerprint TEXT NOT NULL DEFAULT '', + log_summary TEXT NOT NULL DEFAULT '', + metadata_json TEXT NOT NULL DEFAULT '{}', + error TEXT NOT NULL DEFAULT '', + started_at INTEGER NOT NULL DEFAULT 0, + ended_at INTEGER NOT NULL DEFAULT 0, + duration_ms INTEGER NOT NULL DEFAULT 0, + created_by TEXT NOT NULL DEFAULT '', + created_at INTEGER NOT NULL, + updated_at INTEGER NOT NULL +); + CREATE TABLE IF NOT EXISTS plugin_config_snapshots ( id INTEGER PRIMARY KEY AUTOINCREMENT, plugin_id TEXT NOT NULL, @@ -161,6 +199,8 @@ CREATE INDEX IF NOT EXISTS idx_audit_logs_created_at ON audit_logs(created_at); CREATE INDEX IF NOT EXISTS idx_plugin_artifacts_plugin_id ON plugin_artifacts(plugin_id, created_at); CREATE INDEX IF NOT EXISTS idx_plugins_desired_state ON plugins(desired_state, priority); CREATE INDEX IF NOT EXISTS idx_plugin_operations_plugin_id ON plugin_operations(plugin_id, created_at); +CREATE INDEX IF NOT EXISTS idx_plugin_builds_plugin_id ON plugin_builds(plugin_id, created_at); +CREATE INDEX IF NOT EXISTS idx_plugin_builds_source_id ON plugin_builds(source_id, created_at); CREATE INDEX IF NOT EXISTS idx_plugin_config_snapshots_plugin_id ON plugin_config_snapshots(plugin_id, created_at); INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES (1, strftime('%s','now')); ` diff --git a/internal/adminhttp/api.go b/internal/adminhttp/api.go index 1bcb8cf..7dc3f93 100644 --- a/internal/adminhttp/api.go +++ b/internal/adminhttp/api.go @@ -31,6 +31,10 @@ type APIHandlers struct { PluginArtifacts http.HandlerFunc PluginArtifact SegmentHandlerFunc + PluginSources http.HandlerFunc + PluginBuilds http.HandlerFunc + PluginBuild SegmentHandlerFunc + PluginGC http.HandlerFunc PluginsList http.HandlerFunc PluginItem SegmentHandlerFunc PluginAction SegmentHandlerFunc @@ -81,6 +85,14 @@ func NewAPIHandler(prefix string, handlers APIHandlers) http.HandlerFunc { callHandler(w, r, handlers.PluginArtifacts) case strings.HasPrefix(path, "/plugin-artifacts/"): callSegmentHandler(w, r, handlers.PluginArtifact, strings.TrimPrefix(path, "/plugin-artifacts/")) + case path == "/plugin-sources" && (r.Method == http.MethodGet || r.Method == http.MethodPost): + callHandler(w, r, handlers.PluginSources) + case path == "/plugin-builds" && (r.Method == http.MethodGet || r.Method == http.MethodPost): + callHandler(w, r, handlers.PluginBuilds) + case strings.HasPrefix(path, "/plugin-builds/"): + callSegmentHandler(w, r, handlers.PluginBuild, strings.TrimPrefix(path, "/plugin-builds/")) + case path == "/plugin-gc" && (r.Method == http.MethodGet || r.Method == http.MethodPost): + callHandler(w, r, handlers.PluginGC) case path == "/plugins" && r.Method == http.MethodGet: callHandler(w, r, handlers.PluginsList) case path == "/plugins/dispatch-plan" && r.Method == http.MethodGet: diff --git a/internal/adminhttp/api_test.go b/internal/adminhttp/api_test.go index 5b75deb..85cbd39 100644 --- a/internal/adminhttp/api_test.go +++ b/internal/adminhttp/api_test.go @@ -32,6 +32,11 @@ func TestNewAPIHandlerRoutesRequests(t *testing.T) { {name: "audit logs", method: http.MethodGet, path: "/admin/api/audit-logs", wantCall: "audit_logs"}, {name: "plugin artifacts", method: http.MethodGet, path: "/admin/api/plugin-artifacts", wantCall: "plugin_artifacts"}, {name: "plugin artifact", method: http.MethodGet, path: "/admin/api/plugin-artifacts/abc", wantCall: "plugin_artifact", wantSegment: "abc"}, + {name: "plugin sources", method: http.MethodGet, path: "/admin/api/plugin-sources", wantCall: "plugin_sources"}, + {name: "plugin builds", method: http.MethodGet, path: "/admin/api/plugin-builds", wantCall: "plugin_builds"}, + {name: "plugin build", method: http.MethodGet, path: "/admin/api/plugin-builds/7", wantCall: "plugin_build", wantSegment: "7"}, + {name: "plugin build retry", method: http.MethodPost, path: "/admin/api/plugin-builds/7/retry", wantCall: "plugin_build", wantSegment: "7/retry"}, + {name: "plugin gc", method: http.MethodGet, path: "/admin/api/plugin-gc", wantCall: "plugin_gc"}, {name: "plugins list", method: http.MethodGet, path: "/admin/api/plugins", wantCall: "plugins_list"}, {name: "plugin item", method: http.MethodPut, path: "/admin/api/plugins/upstream-rewrite", wantCall: "plugin_item", wantSegment: "upstream-rewrite"}, {name: "plugin action", method: http.MethodPost, path: "/admin/api/plugins/upstream-rewrite/enable", wantCall: "plugin_action", wantSegment: "upstream-rewrite/enable"}, @@ -66,6 +71,10 @@ func TestNewAPIHandlerRoutesRequests(t *testing.T) { PluginArtifacts: recordCall(&gotCall, "plugin_artifacts"), PluginArtifact: recordSegmentCall(&gotCall, &gotSegment, "plugin_artifact"), + PluginSources: recordCall(&gotCall, "plugin_sources"), + PluginBuilds: recordCall(&gotCall, "plugin_builds"), + PluginBuild: recordSegmentCall(&gotCall, &gotSegment, "plugin_build"), + PluginGC: recordCall(&gotCall, "plugin_gc"), PluginsList: recordCall(&gotCall, "plugins_list"), PluginItem: recordSegmentCall(&gotCall, &gotSegment, "plugin_item"), PluginAction: recordSegmentCall(&gotCall, &gotSegment, "plugin_action"), diff --git a/internal/pluginmanager/artifact.go b/internal/pluginmanager/artifact.go index a25a9fa..6b59190 100644 --- a/internal/pluginmanager/artifact.go +++ b/internal/pluginmanager/artifact.go @@ -49,6 +49,94 @@ func NewArtifactStore(root string) ArtifactStore { } func (s ArtifactStore) ValidateAndStore(upload ArtifactUpload) (ArtifactRecord, error) { + return s.validateAndStore(upload, "") +} + +func (s ArtifactStore) ValidateAndStoreBinary(upload ArtifactUpload) (ArtifactRecord, error) { + return s.validateAndStore(upload, ArtifactTypeBinary) +} + +func (s ArtifactStore) ValidateAndStoreSource(upload ArtifactUpload) (ArtifactRecord, error) { + return s.validateAndStore(upload, ArtifactTypeSource) +} + +func (s ArtifactStore) StoreBuiltBinary(upload ArtifactUpload, manifest Manifest, pluginBytes []byte, packageSHA string, metadata map[string]any) (ArtifactRecord, error) { + if s.Root == "" { + return ArtifactRecord{}, errors.New("plugin artifact root is empty") + } + if s.now == nil { + s.now = time.Now + } + manifest.ArtifactType = ArtifactTypeBinary + manifest.Runtime.Entry = RuntimeEntry + if err := validateManifest(manifest); err != nil { + return ArtifactRecord{}, err + } + if len(pluginBytes) == 0 { + return ArtifactRecord{}, errors.New("built runtime entry is empty") + } + pluginSum := sha256.Sum256(pluginBytes) + artifactID := hex.EncodeToString(pluginSum[:]) + artifactDir := filepath.Join(s.Root, manifest.ID, artifactID) + if err := os.MkdirAll(artifactDir, 0755); err != nil { + return ArtifactRecord{}, err + } + pluginPath := filepath.Join(artifactDir, RuntimeEntry) + if err := os.WriteFile(pluginPath, pluginBytes, 0644); err != nil { + return ArtifactRecord{}, err + } + manifestBytes, err := json.MarshalIndent(manifest, "", " ") + if err != nil { + return ArtifactRecord{}, err + } + if err := os.WriteFile(filepath.Join(artifactDir, "manifest.json"), manifestBytes, 0644); err != nil { + return ArtifactRecord{}, err + } + if len(metadata) > 0 { + provenance, err := json.MarshalIndent(metadata, "", " ") + if err != nil { + return ArtifactRecord{}, err + } + if err := os.WriteFile(filepath.Join(artifactDir, "provenance.json"), provenance, 0644); err != nil { + return ArtifactRecord{}, err + } + } + extensionPoints, err := json.Marshal(extensionPointKeys(manifest)) + if err != nil { + return ArtifactRecord{}, err + } + capabilities, err := capabilitiesSummaryJSON(manifest.Capabilities) + if err != nil { + return ArtifactRecord{}, err + } + now := s.now().Unix() + return ArtifactRecord{ + ID: artifactID, + PluginID: manifest.ID, + Version: manifest.Version, + FileName: upload.FileName, + FilePath: pluginPath, + SHA256: artifactID, + PackageSHA256: packageSHA, + SizeBytes: int64(len(pluginBytes)), + ArtifactType: ArtifactTypeBinary, + RuntimeType: manifest.Runtime.Type, + RuntimeEntry: manifest.Runtime.Entry, + Status: ArtifactStatusLoadable, + MetadataJSON: string(manifestBytes), + CapabilitiesSummaryJSON: string(capabilities), + ExtensionPointsJSON: string(extensionPoints), + APIVersion: manifest.APIVersion, + GoVersion: manifest.GoVersion, + GOOS: manifest.GOOS, + GOARCH: manifest.GOARCH, + UploadedBy: upload.Actor, + CreatedAt: now, + UpdatedAt: now, + }, nil +} + +func (s ArtifactStore) validateAndStore(upload ArtifactUpload, expectedArtifactType string) (ArtifactRecord, error) { if s.Root == "" { return ArtifactRecord{}, errors.New("plugin artifact root is empty") } @@ -100,13 +188,16 @@ func (s ArtifactStore) ValidateAndStore(upload ArtifactUpload) (ArtifactRecord, entries := make(map[string]*zip.File) var extractedSize uint64 for _, file := range reader.File { + if file.FileInfo().IsDir() { + if _, err := cleanZipDirName(file.Name); err != nil { + return ArtifactRecord{}, err + } + continue + } clean, err := cleanZipName(file.Name) if err != nil { return ArtifactRecord{}, err } - if file.FileInfo().IsDir() { - continue - } mode := file.FileInfo().Mode() if !mode.IsRegular() || mode&os.ModeType != 0 { return ArtifactRecord{}, fmt.Errorf("unsupported zip entry type %q", file.Name) @@ -141,6 +232,12 @@ func (s ArtifactStore) ValidateAndStore(upload ArtifactUpload) (ArtifactRecord, if err := validateManifest(manifest); err != nil { return ArtifactRecord{}, err } + if expectedArtifactType != "" && manifest.ArtifactType != expectedArtifactType { + return ArtifactRecord{}, fmt.Errorf("artifact_type %q does not match expected %q", manifest.ArtifactType, expectedArtifactType) + } + if manifest.ArtifactType == ArtifactTypeSource { + return s.storeSourcePackage(upload, manifest, manifestBytes, entries, packageSHA) + } entry := manifest.Runtime.Entry pluginFile, ok := entries[entry] @@ -219,6 +316,82 @@ func (s ArtifactStore) ValidateAndStore(upload ArtifactUpload) (ArtifactRecord, }, nil } +func (s ArtifactStore) storeSourcePackage(upload ArtifactUpload, manifest Manifest, manifestBytes []byte, entries map[string]*zip.File, packageSHA string) (ArtifactRecord, error) { + if err := validateSourceEntries(manifest, entries); err != nil { + return ArtifactRecord{}, err + } + artifactID := packageSHA + artifactDir := filepath.Join(s.Root, manifest.ID, artifactID) + sourceDir := filepath.Join(artifactDir, "source") + if err := os.MkdirAll(sourceDir, 0755); err != nil { + return ArtifactRecord{}, err + } + if err := copyFile(upload.SourcePath, filepath.Join(artifactDir, "source.mcgp")); err != nil { + return ArtifactRecord{}, err + } + var sizeBytes int64 + for name, file := range entries { + if name == "manifest.json" { + continue + } + if file.UncompressedSize64 > uint64(s.MaxNonRuntimeBytes) && !strings.HasPrefix(name, "vendor/") { + return ArtifactRecord{}, fmt.Errorf("zip entry %q size %d exceeds limit %d", name, file.UncompressedSize64, s.MaxNonRuntimeBytes) + } + target := filepath.Join(sourceDir, filepath.FromSlash(name)) + if err := os.MkdirAll(filepath.Dir(target), 0755); err != nil { + return ArtifactRecord{}, err + } + data, err := readZipFile(file, s.MaxExtractedBytes) + if err != nil { + return ArtifactRecord{}, err + } + sizeBytes += int64(len(data)) + if err := os.WriteFile(target, data, 0644); err != nil { + return ArtifactRecord{}, err + } + } + if err := os.WriteFile(filepath.Join(artifactDir, "manifest.json"), manifestBytes, 0644); err != nil { + return ArtifactRecord{}, err + } + metadataJSON, err := json.Marshal(manifest) + if err != nil { + return ArtifactRecord{}, err + } + extensionPoints, err := json.Marshal(extensionPointKeys(manifest)) + if err != nil { + return ArtifactRecord{}, err + } + capabilities, err := capabilitiesSummaryJSON(manifest.Capabilities) + if err != nil { + return ArtifactRecord{}, err + } + now := s.now().Unix() + return ArtifactRecord{ + ID: artifactID, + PluginID: manifest.ID, + Version: manifest.Version, + FileName: upload.FileName, + FilePath: sourceDir, + SHA256: artifactID, + PackageSHA256: packageSHA, + SizeBytes: sizeBytes, + ArtifactType: ArtifactTypeSource, + RuntimeType: manifest.Runtime.Type, + RuntimeEntry: sourceBuildEntry(manifest), + Status: ArtifactStatusValidated, + MetadataJSON: string(metadataJSON), + CapabilitiesSummaryJSON: string(capabilities), + ExtensionPointsJSON: string(extensionPoints), + APIVersion: manifest.APIVersion, + GoVersion: manifest.GoVersion, + GOOS: manifest.GOOS, + GOARCH: manifest.GOARCH, + UploadedBy: upload.Actor, + CreatedAt: now, + UpdatedAt: now, + }, nil +} + func capabilitiesSummaryJSON(raw json.RawMessage) ([]byte, error) { summary := CapabilitySummary{ UpstreamConnect: UpstreamConnectCapability{Mode: UpstreamModeDialer}, @@ -259,25 +432,27 @@ func validateManifest(manifest Manifest) error { return fmt.Errorf("invalid plugin id %q", manifest.ID) case strings.TrimSpace(manifest.Version) == "": return errors.New("version is required") - case manifest.ArtifactType != ArtifactTypeBinary: + case manifest.ArtifactType != ArtifactTypeBinary && manifest.ArtifactType != ArtifactTypeSource: return fmt.Errorf("unsupported artifact_type %q", manifest.ArtifactType) case manifest.Runtime.Type != RuntimeGoPlugin: return fmt.Errorf("unsupported runtime.type %q", manifest.Runtime.Type) - case manifest.Runtime.Entry != RuntimeEntry: + case manifest.ArtifactType == ArtifactTypeBinary && manifest.Runtime.Entry != RuntimeEntry: return fmt.Errorf("unsupported runtime.entry %q", manifest.Runtime.Entry) + case manifest.ArtifactType == ArtifactTypeSource && rawSourceBuildEntry(manifest) == "": + return errors.New("build.entry is required for source artifacts") case manifest.APIVersion != APIVersion: return fmt.Errorf("unsupported api_version %q", manifest.APIVersion) - case manifest.GoVersion == "": + case manifest.ArtifactType == ArtifactTypeBinary && manifest.GoVersion == "": return errors.New("go_version is required") - case manifest.GOOS == "": + case manifest.ArtifactType == ArtifactTypeBinary && manifest.GOOS == "": return errors.New("go_os is required") - case manifest.GOARCH == "": + case manifest.ArtifactType == ArtifactTypeBinary && manifest.GOARCH == "": return errors.New("go_arch is required") } - if manifest.GOOS != runtime.GOOS { + if manifest.ArtifactType == ArtifactTypeBinary && manifest.GOOS != "" && manifest.GOOS != runtime.GOOS { return fmt.Errorf("go_os %q does not match gateway %q", manifest.GOOS, runtime.GOOS) } - if manifest.GOARCH != runtime.GOARCH { + if manifest.ArtifactType == ArtifactTypeBinary && manifest.GOARCH != "" && manifest.GOARCH != runtime.GOARCH { return fmt.Errorf("go_arch %q does not match gateway %q", manifest.GOARCH, runtime.GOARCH) } found := false @@ -292,6 +467,65 @@ func validateManifest(manifest Manifest) error { return nil } +func validateSourceEntries(manifest Manifest, entries map[string]*zip.File) error { + if _, ok := entries["go.mod"]; !ok { + return errors.New("source package requires go.mod") + } + buildEntry := sourceBuildEntry(manifest) + if buildEntry == "" || buildEntry == "." { + buildEntry = "." + } + cleanBuildEntry, err := cleanZipName(buildEntry) + if buildEntry == "." { + cleanBuildEntry = "." + err = nil + } + if err != nil { + return fmt.Errorf("invalid runtime.build_entry: %w", err) + } + hasBuildSource := false + hasAnySource := false + for name := range entries { + switch { + case name == "manifest.json" || name == "go.mod" || name == "go.sum": + case strings.HasPrefix(name, "vendor/"): + case strings.EqualFold(path.Base(name), "README.md"), strings.EqualFold(path.Base(name), "LICENSE"), strings.Contains(strings.ToLower(path.Base(name)), "sbom"): + case strings.HasSuffix(name, ".go"): + default: + return fmt.Errorf("unsupported source package entry %q", name) + } + if strings.HasSuffix(name, ".go") { + hasAnySource = true + if cleanBuildEntry == "." || strings.HasPrefix(name, cleanBuildEntry+"/") || path.Dir(name) == cleanBuildEntry { + hasBuildSource = true + } + } + } + if !hasAnySource { + return errors.New("source package requires at least one Go source file") + } + if !hasBuildSource { + return fmt.Errorf("source package build entry %q has no Go source files", manifest.Runtime.BuildEntry) + } + return nil +} + +func sourceBuildEntry(manifest Manifest) string { + entry := rawSourceBuildEntry(manifest) + if entry == "" { + return SourceBuildEntry + } + return entry +} + +func rawSourceBuildEntry(manifest Manifest) string { + entry := strings.Trim(strings.TrimSpace(manifest.Build.Entry), "/") + if entry == "" { + entry = strings.Trim(strings.TrimSpace(manifest.Runtime.BuildEntry), "/") + } + return entry +} + func extensionPointKeys(manifest Manifest) []string { keys := make([]string, 0, len(manifest.ExtensionPoints)) for _, ep := range manifest.ExtensionPoints { @@ -311,6 +545,14 @@ func cleanZipName(name string) (string, error) { return clean, nil } +func cleanZipDirName(name string) (string, error) { + name = strings.TrimSuffix(name, "/") + if name == "" { + return "", fmt.Errorf("unsafe zip entry %q", name) + } + return cleanZipName(name) +} + func readZipFile(file *zip.File, maxBytes int64) ([]byte, error) { rc, err := file.Open() if err != nil { @@ -341,3 +583,21 @@ func fileSHA256(path string) (string, error) { } return hex.EncodeToString(hash.Sum(nil)), nil } + +func copyFile(src, dst string) error { + input, err := os.Open(src) + if err != nil { + return err + } + defer input.Close() + + output, err := os.Create(dst) + if err != nil { + return err + } + defer output.Close() + if _, err := io.Copy(output, input); err != nil { + return err + } + return output.Close() +} diff --git a/internal/pluginmanager/artifact_test.go b/internal/pluginmanager/artifact_test.go index 56b59cb..9f98b1f 100644 --- a/internal/pluginmanager/artifact_test.go +++ b/internal/pluginmanager/artifact_test.go @@ -38,6 +38,30 @@ func TestArtifactStoreValidateAndStore(t *testing.T) { } } +func TestArtifactStoreValidateAndStoreSource(t *testing.T) { + packagePath := writeTestMCGP(t, map[string][]byte{ + "manifest.json": testSourceManifestBytes(t, "source-plugin"), + "go.mod": []byte("module example.com/source-plugin\n\ngo 1.24.0\n"), + "main.go": []byte("package main\n"), + "README.md": []byte("source fixture"), + }) + store := NewArtifactStore(t.TempDir()) + source, err := store.ValidateAndStoreSource(ArtifactUpload{ + SourcePath: packagePath, + FileName: "source-plugin.mcgp", + Actor: "admin", + }) + if err != nil { + t.Fatalf("ValidateAndStoreSource() error = %v", err) + } + if source.ArtifactType != ArtifactTypeSource || source.Status != ArtifactStatusValidated { + t.Fatalf("source = %+v, want validated source", source) + } + if _, err := os.Stat(filepath.Join(source.FilePath, "go.mod")); err != nil { + t.Fatalf("stored source go.mod stat error = %v", err) + } +} + func TestArtifactStoreRejectsUnsafePackage(t *testing.T) { tests := []struct { name string @@ -90,6 +114,60 @@ func TestArtifactStoreRejectsUnsafePackage(t *testing.T) { } } +func TestArtifactStoreRejectsSourceShellScripts(t *testing.T) { + store := NewArtifactStore(t.TempDir()) + _, err := store.ValidateAndStoreSource(ArtifactUpload{ + SourcePath: writeTestMCGP(t, map[string][]byte{ + "manifest.json": testSourceManifestBytes(t, "test-plugin"), + "go.mod": []byte("module example.com/test\n"), + "main.go": []byte("package main\n"), + "build.sh": []byte("go build"), + }), + FileName: "bad-source.mcgp", + Actor: "admin", + }) + if err == nil || !strings.Contains(err.Error(), "unsupported source package entry") { + t.Fatalf("ValidateAndStoreSource() error = %v, want unsupported source entry", err) + } +} + +func testSourceManifestBytes(t *testing.T, pluginID string) []byte { + t.Helper() + manifest := Manifest{ + SchemaVersion: SchemaVersion, + ID: pluginID, + Name: "Source Plugin", + Version: "0.1.0", + ArtifactType: ArtifactTypeSource, + Runtime: RuntimeManifest{ + Type: RuntimeGoPlugin, + EntrySymbol: "Plugin", + }, + Build: BuildManifest{ + Type: BuildTypeGo, + Entry: ".", + GoVersion: runtime.Version(), + Tags: []string{}, + VendorRequired: false, + Output: RuntimeEntry, + }, + APIVersion: APIVersion, + GoVersion: runtime.Version(), + GOOS: runtime.GOOS, + GOARCH: runtime.GOARCH, + ExtensionPoints: []ExtensionPoint{{ + Type: "hook", + Key: ExtensionUpstreamConnect, + }}, + Capabilities: json.RawMessage(`{"extension_points":["upstream.connect/v1"]}`), + } + data, err := json.Marshal(manifest) + if err != nil { + t.Fatalf("Marshal source manifest error = %v", err) + } + return data +} + func testManifestBytes(t *testing.T, pluginID string) []byte { return testManifestBytesWithCapabilities(t, pluginID, json.RawMessage(`{"extension_points":["upstream.connect/v1"]}`)) } diff --git a/internal/pluginmanager/builder.go b/internal/pluginmanager/builder.go new file mode 100644 index 0000000..1ee41ef --- /dev/null +++ b/internal/pluginmanager/builder.go @@ -0,0 +1,336 @@ +package pluginmanager + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "runtime" + "strings" + "time" +) + +type SourceBuilder interface { + Build(ctx context.Context, source ArtifactRecord, req BuildRequest, build BuildRecord) (BuildResult, error) +} + +type BuildResult struct { + Manifest Manifest + ArtifactBytes []byte + ArtifactSHA256 string + GoVersion string + ModuleSummary string + GoVersionM string + ABIFingerprint string + LogSummary string + Metadata map[string]any +} + +type LocalProcessBuilder struct { + StoreRoot string +} + +func (b LocalProcessBuilder) Build(ctx context.Context, source ArtifactRecord, req BuildRequest, build BuildRecord) (BuildResult, error) { + if source.ArtifactType != ArtifactTypeSource { + return BuildResult{}, fmt.Errorf("artifact %s is %q, want source", source.ID, source.ArtifactType) + } + var manifest Manifest + if err := json.Unmarshal([]byte(source.MetadataJSON), &manifest); err != nil { + return BuildResult{}, fmt.Errorf("decode source manifest: %w", err) + } + if err := validateManifest(manifest); err != nil { + return BuildResult{}, err + } + if manifest.Build.Type != "" && manifest.Build.Type != BuildTypeGo { + return BuildResult{}, fmt.Errorf("unsupported build.type %q", manifest.Build.Type) + } + if manifest.Build.Output != "" && manifest.Build.Output != RuntimeEntry { + return BuildResult{}, fmt.Errorf("unsupported build.output %q", manifest.Build.Output) + } + if req.VendorRequired { + if _, err := os.Stat(filepath.Join(source.FilePath, "vendor")); err != nil { + return BuildResult{}, errors.New("vendor is required but source package has no vendor directory") + } + } + goVersion, goVersionOutput, err := commandOutput(ctx, source.FilePath, nil, "go", "version") + if err != nil { + return BuildResult{LogSummary: sanitizeLog(goVersionOutput)}, err + } + if manifest.GoVersion != "" && manifest.GoVersion != goVersion { + return BuildResult{GoVersion: goVersion, LogSummary: sanitizeLog(goVersionOutput)}, fmt.Errorf("source go_version %q does not match builder %q", manifest.GoVersion, goVersion) + } + manifest.GoVersion = goVersion + manifest.GOOS = req.GOOS + manifest.GOARCH = req.GOARCH + manifest.ArtifactType = ArtifactTypeBinary + manifest.Runtime.Entry = RuntimeEntry + + outDir, err := os.MkdirTemp("", "mc-gateway-plugin-build-*") + if err != nil { + return BuildResult{}, err + } + defer os.RemoveAll(outDir) + outPath := filepath.Join(outDir, RuntimeEntry) + env := buildEnvironment(req) + args := []string{"build", "-buildmode=plugin", "-trimpath", "-buildvcs=false", "-o", outPath} + if req.BuildTags != "" { + args = append(args, "-tags", req.BuildTags) + } + buildEntry := sourceBuildEntry(manifest) + args = append(args, buildEntry) + _, buildLog, buildErr := commandOutput(ctx, source.FilePath, env, "go", args...) + if buildErr != nil { + return BuildResult{GoVersion: goVersion, LogSummary: sanitizeLog(buildLog)}, buildErr + } + if _, nmLog, err := commandOutput(ctx, source.FilePath, env, "go", "tool", "nm", outPath); err != nil { + return BuildResult{GoVersion: goVersion, LogSummary: sanitizeLog(buildLog + "\n" + nmLog)}, fmt.Errorf("inspect built plugin symbols: %w", err) + } else if err := validateBuiltSymbols(manifest, nmLog); err != nil { + return BuildResult{GoVersion: goVersion, LogSummary: sanitizeLog(buildLog + "\n" + nmLog)}, err + } + artifactBytes, err := os.ReadFile(outPath) + if err != nil { + return BuildResult{}, err + } + artifactSum := sha256.Sum256(artifactBytes) + artifactSHA := hex.EncodeToString(artifactSum[:]) + moduleSummary := "[]" + if _, moduleLog, err := commandOutput(ctx, source.FilePath, env, "go", "list", "-m", "-json", "all"); err == nil { + moduleSummary = summarizeGoModules(moduleLog) + } else { + buildLog += "\n" + moduleLog + } + goVersionM := "{}" + if _, versionMLog, err := commandOutput(ctx, source.FilePath, env, "go", "version", "-m", outPath); err == nil { + goVersionM = summarizeGoVersionM(versionMLog) + } else { + buildLog += "\n" + versionMLog + } + abi := abiFingerprint(manifest, goVersion) + metadata := map[string]any{ + "source_sha256": source.SHA256, + "artifact_sha256": artifactSHA, + "builder_type": BuilderTypeLocalProcess, + "builder_version": build.BuilderVersion, + "go_version": goVersion, + "go_os": req.GOOS, + "go_arch": req.GOARCH, + "go_amd64": req.GOAMD64, + "go_arm64": req.GOARM64, + "cgo_enabled": req.CGOEnabled, + "build_tags": req.BuildTags, + "vendor_required": req.VendorRequired, + "abi_fingerprint": abi, + } + return BuildResult{ + Manifest: manifest, + ArtifactBytes: artifactBytes, + ArtifactSHA256: artifactSHA, + GoVersion: goVersion, + ModuleSummary: moduleSummary, + GoVersionM: goVersionM, + ABIFingerprint: abi, + LogSummary: sanitizeLog(buildLog), + Metadata: metadata, + }, nil +} + +type ContainerBuilder struct{} + +func (ContainerBuilder) Build(context.Context, ArtifactRecord, BuildRequest, BuildRecord) (BuildResult, error) { + return BuildResult{}, errors.New("container builder is not configured in this phase") +} + +func buildEnvironment(req BuildRequest) []string { + env := os.Environ() + env = append(env, + "GOOS="+req.GOOS, + "GOARCH="+req.GOARCH, + "CGO_ENABLED="+req.CGOEnabled, + ) + if req.GOAMD64 != "" { + env = append(env, "GOAMD64="+req.GOAMD64) + } + if req.GOARM64 != "" { + env = append(env, "GOARM64="+req.GOARM64) + } + if req.GOPROXY != "" { + env = append(env, "GOPROXY="+req.GOPROXY) + } + if req.GONOSUMDB != "" { + env = append(env, "GONOSUMDB="+req.GONOSUMDB) + } + if req.GOPRIVATE != "" { + env = append(env, "GOPRIVATE="+req.GOPRIVATE) + } + return env +} + +func commandOutput(ctx context.Context, dir string, env []string, name string, args ...string) (string, string, error) { + cmd := exec.CommandContext(ctx, name, args...) + cmd.Dir = dir + if env != nil { + cmd.Env = env + } + var out bytes.Buffer + cmd.Stdout = &out + cmd.Stderr = &out + err := cmd.Run() + output := out.String() + if name == "go" && len(args) == 1 && args[0] == "version" && err == nil { + fields := strings.Fields(output) + if len(fields) >= 3 { + return fields[2], output, nil + } + } + return strings.TrimSpace(output), output, err +} + +func sanitizeLog(log string) string { + replacers := []string{ + os.Getenv("HOME"), "$HOME", + os.TempDir(), "$TMPDIR", + } + sanitized := log + for i := 0; i+1 < len(replacers); i += 2 { + if replacers[i] != "" { + sanitized = strings.ReplaceAll(sanitized, replacers[i], replacers[i+1]) + } + } + for _, key := range []string{"TOKEN", "SECRET", "PASSWORD", "PRIVATE"} { + for _, part := range strings.Fields(sanitized) { + if strings.Contains(strings.ToUpper(part), key+"=") { + sanitized = strings.ReplaceAll(sanitized, part, key+"=") + } + } + } + if len(sanitized) > DefaultBuildLogMaxBytes { + sanitized = sanitized[len(sanitized)-DefaultBuildLogMaxBytes:] + } + lines := strings.Split(sanitized, "\n") + if len(lines) > 80 { + lines = lines[len(lines)-80:] + } + return strings.TrimSpace(strings.Join(lines, "\n")) +} + +func summarizeGoModules(raw string) string { + decoder := json.NewDecoder(strings.NewReader(raw)) + var modules []map[string]any + for decoder.More() { + var module struct { + Path string `json:"Path"` + Version string `json:"Version"` + Main bool `json:"Main"` + Replace *struct { + Path string `json:"Path"` + Version string `json:"Version"` + } `json:"Replace"` + } + if err := decoder.Decode(&module); err != nil { + break + } + item := map[string]any{"path": module.Path, "version": module.Version, "main": module.Main} + if module.Replace != nil { + item["replace"] = map[string]string{"path": module.Replace.Path, "version": module.Replace.Version} + } + modules = append(modules, item) + } + data, err := json.Marshal(modules) + if err != nil || len(data) == 0 { + return "[]" + } + return string(data) +} + +func summarizeGoVersionM(raw string) string { + summary := map[string]any{"raw_summary": sanitizeLog(raw)} + data, err := json.Marshal(summary) + if err != nil { + return "{}" + } + return string(data) +} + +func abiFingerprint(manifest Manifest, goVersion string) string { + payload := strings.Join([]string{ + manifest.ID, + manifest.APIVersion, + manifest.SDKModule, + manifest.SDKModuleVersion, + goVersion, + runtime.GOOS, + runtime.GOARCH, + extensionPointsFingerprint(manifest), + }, "\x00") + sum := sha256.Sum256([]byte(payload)) + return hex.EncodeToString(sum[:]) +} + +func extensionPointsFingerprint(manifest Manifest) string { + keys := extensionPointKeys(manifest) + data, _ := json.Marshal(keys) + return string(data) +} + +func validateBuiltSymbols(manifest Manifest, nmLog string) error { + entrySymbol := manifest.Runtime.EntrySymbol + if entrySymbol == "" { + entrySymbol = "Plugin" + } + metadataSymbol := manifest.Runtime.MetadataSymbol + if metadataSymbol == "" { + metadataSymbol = "MCGatewayPluginMetadata" + } + if !strings.Contains(nmLog, entrySymbol) { + return fmt.Errorf("built plugin is missing entry symbol %q", entrySymbol) + } + if !strings.Contains(nmLog, metadataSymbol) { + return fmt.Errorf("built plugin is missing metadata symbol %q", metadataSymbol) + } + return nil +} + +func defaultBuildRequest(req BuildRequest, source ArtifactRecord) BuildRequest { + var manifest Manifest + _ = json.Unmarshal([]byte(source.MetadataJSON), &manifest) + if req.BuilderType == "" { + req.BuilderType = BuilderTypeLocalProcess + } + if req.BuilderVersion == "" { + req.BuilderVersion = "local-process/go-buildmode-plugin" + } + if req.GOOS == "" { + req.GOOS = runtime.GOOS + } + if req.GOARCH == "" { + req.GOARCH = runtime.GOARCH + } + if req.CGOEnabled == "" { + if manifest.Build.CGOEnabled != nil && !*manifest.Build.CGOEnabled { + req.CGOEnabled = "0" + } else { + req.CGOEnabled = "1" + } + } + if req.BuildTags == "" && len(manifest.Build.Tags) > 0 { + req.BuildTags = strings.Join(manifest.Build.Tags, ",") + } + if manifest.Build.VendorRequired { + req.VendorRequired = true + } + if req.SDKModule == "" { + req.SDKModule = manifest.SDKModule + req.SDKVersion = manifest.SDKModuleVersion + } + return req +} + +func buildDurationMS(start time.Time) int64 { + return time.Since(start).Milliseconds() +} diff --git a/internal/pluginmanager/gc.go b/internal/pluginmanager/gc.go new file mode 100644 index 0000000..9894af2 --- /dev/null +++ b/internal/pluginmanager/gc.go @@ -0,0 +1,129 @@ +package pluginmanager + +import ( + "context" + "os" + "path/filepath" + "strconv" +) + +func (m *Manager) GCCandidates(ctx context.Context) ([]GCCandidate, error) { + refs, err := m.repo.ReferencedArtifactIDs(ctx) + if err != nil { + return nil, err + } + artifacts, err := m.repo.ListArtifacts(ctx, "") + if err != nil { + return nil, err + } + var candidates []GCCandidate + for _, artifact := range artifacts { + referenced := refs[artifact.ID] + artifactPath := artifactGCPath(artifact) + size := artifact.SizeBytes + if info, err := os.Stat(artifactPath); err == nil && info.IsDir() { + size = dirSize(artifactPath) + } else if err == nil { + size = info.Size() + } + candidate := GCCandidate{ + Kind: "artifact", + ID: artifact.ID, + PluginID: artifact.PluginID, + Path: artifactPath, + Protected: referenced, + SizeBytes: size, + CreatedAt: artifact.CreatedAt, + Referenced: referenced, + } + switch { + case referenced: + candidate.Reason = "referenced by active/desired/snapshot/build" + case artifact.ArtifactType == ArtifactTypeSource: + candidate.Reason = "unreferenced source package" + case artifact.Status == ArtifactStatusRejected || artifact.Status == ArtifactStatusDeleted: + candidate.Reason = "unreferenced rejected/deleted artifact" + default: + candidate.Reason = "unreferenced artifact" + } + if !referenced && (artifact.ArtifactType == ArtifactTypeSource || artifact.Status == ArtifactStatusRejected || artifact.Status == ArtifactStatusDeleted) { + candidates = append(candidates, candidate) + } + } + builds, err := m.repo.ListBuilds(ctx, "") + if err != nil { + return nil, err + } + for _, build := range builds { + protectedBuild := build.Status == BuildStatusRunning || build.Status == BuildStatusQueued + reason := "completed build log" + if protectedBuild { + reason = "in-flight build" + } + candidates = append(candidates, GCCandidate{ + Kind: "build_log", + ID: buildIDString(build.ID), + PluginID: build.PluginID, + Protected: protectedBuild, + Reason: reason, + SizeBytes: int64(len(build.LogSummary)), + CreatedAt: build.CreatedAt, + }) + } + return candidates, nil +} + +func (m *Manager) RunGC(ctx context.Context, actor string, dryRun bool) ([]GCCandidate, error) { + candidates, err := m.GCCandidates(ctx) + if err != nil { + return nil, err + } + if dryRun { + _ = m.repo.RecordOperation(ctx, "", "", "artifact_gc", "dry_run", actor, "artifact gc dry-run completed", map[string]any{ + "candidates": len(candidates), + }) + return candidates, nil + } + var removed []GCCandidate + for _, candidate := range candidates { + if candidate.Protected || candidate.Path == "" || candidate.Kind != "artifact" { + continue + } + if err := os.RemoveAll(candidate.Path); err != nil { + _ = m.repo.RecordOperation(ctx, candidate.PluginID, candidate.ID, "artifact_gc", "failed", actor, err.Error(), map[string]any{ + "path": candidate.Path, + }) + continue + } + removed = append(removed, candidate) + } + _ = m.repo.RecordOperation(ctx, "", "", "artifact_gc", "succeeded", actor, "artifact gc completed", map[string]any{ + "removed": len(removed), + }) + return removed, nil +} + +func artifactGCPath(artifact ArtifactRecord) string { + if artifact.FilePath == "" { + return "" + } + return filepath.Dir(artifact.FilePath) +} + +func dirSize(root string) int64 { + var total int64 + _ = filepath.WalkDir(root, func(path string, d os.DirEntry, err error) error { + if err != nil || d.IsDir() { + return nil + } + if info, err := d.Info(); err == nil { + total += info.Size() + } + return nil + }) + return total +} + +func buildIDString(id int64) string { + return strconv.FormatInt(id, 10) +} diff --git a/internal/pluginmanager/manager.go b/internal/pluginmanager/manager.go index 4e0826e..2f431eb 100644 --- a/internal/pluginmanager/manager.go +++ b/internal/pluginmanager/manager.go @@ -10,6 +10,7 @@ import ( "net" stdplugin "plugin" "reflect" + "runtime" "sort" "sync" "sync/atomic" @@ -71,6 +72,7 @@ type Manager struct { repo Repository store ArtifactStore adapter RuntimeAdapter + builders map[string]SourceBuilder handleConn func(net.Conn) wg *sync.WaitGroup @@ -147,6 +149,7 @@ type Options struct { HandleConn func(net.Conn) WaitGroup *sync.WaitGroup Adapter RuntimeAdapter + Builders map[string]SourceBuilder } func New(options Options) *Manager { @@ -158,12 +161,19 @@ func New(options Options) *Manager { repo: NewRepository(options.DB), store: NewArtifactStore(options.ArtifactRoot), adapter: adapter, + builders: options.Builders, handleConn: options.HandleConn, wg: options.WaitGroup, loaded: make(map[string]*loadedPlugin), proxyConns: make(map[uint64]*proxyConnection), drainingIDs: make(map[string]bool), } + if manager.builders == nil { + manager.builders = map[string]SourceBuilder{ + BuilderTypeLocalProcess: LocalProcessBuilder{StoreRoot: options.ArtifactRoot}, + BuilderTypeContainer: ContainerBuilder{}, + } + } manager.publish(nil) return manager } @@ -174,6 +184,9 @@ func (m *Manager) UploadArtifact(ctx context.Context, upload ArtifactUpload) (Ar _ = m.repo.RecordOperation(ctx, "", "", "artifact_upload", "failed", upload.Actor, err.Error(), nil) return ArtifactRecord{}, err } + if artifact.ArtifactType == ArtifactTypeSource { + return m.saveSourceArtifact(ctx, upload.Actor, artifact, "artifact_upload") + } if err := m.repo.SaveArtifact(ctx, artifact); err != nil { return ArtifactRecord{}, err } @@ -186,6 +199,203 @@ func (m *Manager) UploadArtifact(ctx context.Context, upload ArtifactUpload) (Ar return artifact, nil } +func (m *Manager) UploadSource(ctx context.Context, upload ArtifactUpload) (ArtifactRecord, error) { + artifact, err := m.store.ValidateAndStoreSource(upload) + if err != nil { + _ = m.repo.RecordOperation(ctx, "", "", "source_upload", "failed", upload.Actor, err.Error(), nil) + return ArtifactRecord{}, err + } + return m.saveSourceArtifact(ctx, upload.Actor, artifact, "source_upload") +} + +func (m *Manager) saveSourceArtifact(ctx context.Context, actor string, artifact ArtifactRecord, operation string) (ArtifactRecord, error) { + if err := m.repo.SaveArtifact(ctx, artifact); err != nil { + return ArtifactRecord{}, err + } + build, err := m.CreateBuild(ctx, actor, BuildRequest{SourceID: artifact.ID}) + if err != nil { + _ = m.repo.RecordOperation(ctx, artifact.PluginID, artifact.ID, "source_build_queue", "failed", actor, err.Error(), nil) + return ArtifactRecord{}, err + } + _ = m.repo.RecordOperation(ctx, artifact.PluginID, artifact.ID, operation, "succeeded", actor, "source package uploaded", map[string]any{ + "source_sha256": artifact.SHA256, + "api_version": artifact.APIVersion, + "go_version": artifact.GoVersion, + "build_id": build.ID, + }) + return artifact, nil +} + +func (m *Manager) BuildSource(ctx context.Context, actor string, req BuildRequest) (BuildRecord, error) { + build, err := m.CreateBuild(ctx, actor, req) + if err != nil { + return BuildRecord{}, err + } + return m.RunBuild(ctx, actor, build.ID) +} + +func (m *Manager) CreateBuild(ctx context.Context, actor string, req BuildRequest) (BuildRecord, error) { + source, err := m.repo.Artifact(ctx, req.SourceID) + if err != nil { + return BuildRecord{}, err + } + if source.ArtifactType != ArtifactTypeSource { + return BuildRecord{}, fmt.Errorf("artifact %s is %q, want source", source.ID, source.ArtifactType) + } + req = defaultBuildRequest(req, source) + if req.GOOS != runtime.GOOS || req.GOARCH != runtime.GOARCH { + return BuildRecord{}, fmt.Errorf("build target %s/%s does not match gateway %s/%s", req.GOOS, req.GOARCH, runtime.GOOS, runtime.GOARCH) + } + builder := m.builders[req.BuilderType] + if builder == nil { + return BuildRecord{}, fmt.Errorf("builder type %q is not available", req.BuilderType) + } + build, err := m.repo.CreateBuild(ctx, BuildRecord{ + PluginID: source.PluginID, + SourceID: source.ID, + Status: BuildStatusQueued, + BuilderType: req.BuilderType, + BuilderImage: req.BuilderImage, + BuilderVersion: req.BuilderVersion, + GOOS: req.GOOS, + GOARCH: req.GOARCH, + GOAMD64: req.GOAMD64, + GOARM64: req.GOARM64, + CGOEnabled: req.CGOEnabled, + BuildTags: req.BuildTags, + SDKModule: req.SDKModule, + SDKVersion: req.SDKVersion, + GOPROXY: req.GOPROXY, + GONOSUMDB: req.GONOSUMDB, + GOPRIVATE: req.GOPRIVATE, + VendorRequired: req.VendorRequired, + SourceSHA256: source.SHA256, + ModuleSummary: "[]", + GoVersionM: "{}", + MetadataJSON: "{}", + CreatedBy: actor, + }) + if err != nil { + return BuildRecord{}, err + } + _ = m.repo.RecordOperation(ctx, source.PluginID, source.ID, "source_build_queue", "succeeded", actor, "source build queued", map[string]any{ + "build_id": build.ID, + "builder_type": req.BuilderType, + "source_sha256": source.SHA256, + }) + return build, nil +} + +func (m *Manager) RunBuild(ctx context.Context, actor string, buildID int64) (BuildRecord, error) { + build, err := m.repo.Build(ctx, buildID) + if err != nil { + return BuildRecord{}, err + } + if build.Status != BuildStatusQueued { + return BuildRecord{}, fmt.Errorf("build status %q cannot be run", build.Status) + } + source, err := m.repo.Artifact(ctx, build.SourceID) + if err != nil { + return BuildRecord{}, err + } + req := BuildRequest{ + SourceID: build.SourceID, + BuilderType: build.BuilderType, + BuilderImage: build.BuilderImage, + BuilderVersion: build.BuilderVersion, + GOOS: build.GOOS, + GOARCH: build.GOARCH, + GOAMD64: build.GOAMD64, + GOARM64: build.GOARM64, + CGOEnabled: build.CGOEnabled, + BuildTags: build.BuildTags, + SDKModule: build.SDKModule, + SDKVersion: build.SDKVersion, + GOPROXY: build.GOPROXY, + GONOSUMDB: build.GONOSUMDB, + GOPRIVATE: build.GOPRIVATE, + VendorRequired: build.VendorRequired, + } + builder := m.builders[build.BuilderType] + if builder == nil { + return BuildRecord{}, fmt.Errorf("builder type %q is not available", build.BuilderType) + } + start := time.Now() + if err := m.repo.MarkBuildRunning(ctx, build.ID); err != nil { + return BuildRecord{}, err + } + build, _ = m.repo.Build(ctx, build.ID) + result, buildErr := builder.Build(ctx, source, req, build) + build.StartedAt = start.Unix() + build.EndedAt = time.Now().Unix() + build.DurationMS = buildDurationMS(start) + build.GoVersion = result.GoVersion + build.ModuleSummary = result.ModuleSummary + build.GoVersionM = result.GoVersionM + build.ABIFingerprint = result.ABIFingerprint + build.LogSummary = result.LogSummary + build.SourceSHA256 = source.SHA256 + if buildErr != nil { + build.Status = BuildStatusFailed + build.Error = buildErr.Error() + _ = m.repo.FinishBuild(ctx, build) + _ = m.repo.RecordOperation(ctx, source.PluginID, source.ID, "source_build", "failed", actor, buildErr.Error(), map[string]any{ + "build_id": build.ID, + "source_sha256": source.SHA256, + "builder_type": req.BuilderType, + "log_summary": result.LogSummary, + "active_changed": false, + }) + return m.repo.Build(ctx, build.ID) + } + if result.Manifest.GoVersion != result.GoVersion { + build.Status = BuildStatusFailed + build.Error = fmt.Sprintf("manifest go_version %q does not match built Go version %q", result.Manifest.GoVersion, result.GoVersion) + _ = m.repo.FinishBuild(ctx, build) + return m.repo.Build(ctx, build.ID) + } + metadata := result.Metadata + if metadata == nil { + metadata = map[string]any{} + } + metadata["build_id"] = build.ID + metadata["module_summary"] = json.RawMessage(defaultJSONArray(result.ModuleSummary)) + metadata["go_version_m"] = json.RawMessage(defaultJSONObject(result.GoVersionM)) + artifact, err := m.store.StoreBuiltBinary(ArtifactUpload{ + FileName: result.Manifest.ID + "-" + result.Manifest.Version + ".mcgp", + Actor: actor, + }, result.Manifest, result.ArtifactBytes, source.PackageSHA256, metadata) + if err != nil { + build.Status = BuildStatusFailed + build.Error = err.Error() + _ = m.repo.FinishBuild(ctx, build) + return m.repo.Build(ctx, build.ID) + } + if err := m.repo.SaveArtifact(ctx, artifact); err != nil { + build.Status = BuildStatusFailed + build.Error = err.Error() + _ = m.repo.FinishBuild(ctx, build) + return m.repo.Build(ctx, build.ID) + } + build.Status = BuildStatusSucceeded + build.ArtifactID = artifact.ID + build.ArtifactSHA256 = artifact.SHA256 + metadataBytes, _ := json.Marshal(metadata) + build.MetadataJSON = string(metadataBytes) + if err := m.repo.FinishBuild(ctx, build); err != nil { + return BuildRecord{}, err + } + _ = m.repo.RecordOperation(ctx, source.PluginID, artifact.ID, "source_build", "succeeded", actor, "source build succeeded", map[string]any{ + "build_id": build.ID, + "source_id": source.ID, + "source_sha256": source.SHA256, + "artifact_sha256": artifact.SHA256, + "builder_type": req.BuilderType, + "go_version": result.GoVersion, + }) + return m.repo.Build(ctx, build.ID) +} + func (m *Manager) SetDesired(ctx context.Context, actor, pluginID, artifactID, desiredState, configJSON string, priority int) (PluginRecord, error) { pluginRecord, err := m.repo.UpsertDesired(ctx, actor, pluginID, artifactID, desiredState, configJSON, priority) if err != nil { @@ -440,6 +650,55 @@ func (m *Manager) ListArtifacts(ctx context.Context, pluginID string) ([]Artifac return m.repo.ListArtifacts(ctx, pluginID) } +func (m *Manager) ListBuilds(ctx context.Context, pluginID string) ([]BuildRecord, error) { + return m.repo.ListBuilds(ctx, pluginID) +} + +func (m *Manager) Build(ctx context.Context, id int64) (BuildRecord, error) { + return m.repo.Build(ctx, id) +} + +func (m *Manager) CancelBuild(ctx context.Context, actor string, id int64) (BuildRecord, error) { + build, err := m.repo.CancelBuild(ctx, id, actor) + if err != nil { + return BuildRecord{}, err + } + _ = m.repo.RecordOperation(ctx, build.PluginID, build.SourceID, "source_build_cancel", "succeeded", actor, "build canceled", map[string]any{"build_id": id}) + return build, nil +} + +func (m *Manager) RetryBuild(ctx context.Context, actor string, id int64) (BuildRecord, error) { + build, err := m.repo.Build(ctx, id) + if err != nil { + return BuildRecord{}, err + } + if build.Status != BuildStatusFailed && build.Status != BuildStatusCanceled { + return BuildRecord{}, fmt.Errorf("build status %q cannot be retried", build.Status) + } + next, err := m.CreateBuild(ctx, actor, BuildRequest{ + SourceID: build.SourceID, + BuilderType: build.BuilderType, + BuilderImage: build.BuilderImage, + BuilderVersion: build.BuilderVersion, + GOOS: build.GOOS, + GOARCH: build.GOARCH, + GOAMD64: build.GOAMD64, + GOARM64: build.GOARM64, + CGOEnabled: build.CGOEnabled, + BuildTags: build.BuildTags, + SDKModule: build.SDKModule, + SDKVersion: build.SDKVersion, + GOPROXY: build.GOPROXY, + GONOSUMDB: build.GONOSUMDB, + GOPRIVATE: build.GOPRIVATE, + VendorRequired: build.VendorRequired, + }) + if err != nil { + return BuildRecord{}, err + } + return m.RunBuild(ctx, actor, next.ID) +} + func (m *Manager) Artifact(ctx context.Context, id string) (ArtifactRecord, error) { return m.repo.Artifact(ctx, id) } @@ -574,6 +833,15 @@ func (m *Manager) loadLocked(ctx context.Context, pluginRecord PluginRecord) (*l if artifact.Status == ArtifactStatusDeleted || artifact.Status == ArtifactStatusRejected { return nil, fmt.Errorf("artifact status %q is not loadable", artifact.Status) } + if artifact.ArtifactType != ArtifactTypeBinary { + return nil, fmt.Errorf("artifact type %q is not loadable", artifact.ArtifactType) + } + if artifact.GoVersion != runtime.Version() { + return nil, fmt.Errorf("artifact go_version %q does not match gateway %q", artifact.GoVersion, runtime.Version()) + } + if artifact.GOOS != runtime.GOOS || artifact.GOARCH != runtime.GOARCH { + return nil, fmt.Errorf("artifact target %s/%s does not match gateway %s/%s", artifact.GOOS, artifact.GOARCH, runtime.GOOS, runtime.GOARCH) + } gateway := NewGateway(pluginRecord.ID, m.handleConn, m.wg) instance, err := m.adapter.Load(ctx, artifact, pluginRecord, gateway) diff --git a/internal/pluginmanager/manager_test.go b/internal/pluginmanager/manager_test.go index b83d0ae..2e0379f 100644 --- a/internal/pluginmanager/manager_test.go +++ b/internal/pluginmanager/manager_test.go @@ -9,6 +9,8 @@ import ( "io" "net" "path/filepath" + "runtime" + "strings" "testing" "time" @@ -29,6 +31,104 @@ func TestManagerUploadDoesNotLoadPlugin(t *testing.T) { } } +func TestManagerUploadSourceQueuesBuild(t *testing.T) { + manager := newManagerForTest(t, &fakeAdapter{}) + source := uploadTestSource(t, manager, "plugin-a") + builds, err := manager.ListBuilds(context.Background(), "plugin-a") + if err != nil { + t.Fatalf("ListBuilds() error = %v", err) + } + if len(builds) != 1 { + t.Fatalf("builds = %d, want 1", len(builds)) + } + if builds[0].SourceID != source.ID || builds[0].Status != BuildStatusQueued { + t.Fatalf("queued build = %+v, want source %s queued", builds[0], source.ID) + } +} + +func TestManagerBuildSourceCreatesBinaryArtifact(t *testing.T) { + t.Setenv("GOCACHE", t.TempDir()) + t.Setenv("GOWORK", "off") + manager := newManagerForTest(t, &fakeAdapter{}) + _ = uploadBuildableTestSource(t, manager, "source-enable") + builds, err := manager.ListBuilds(context.Background(), "source-enable") + if err != nil { + t.Fatalf("ListBuilds() error = %v", err) + } + if len(builds) != 1 { + t.Fatalf("builds = %d, want 1", len(builds)) + } + build, err := manager.RunBuild(context.Background(), "admin", builds[0].ID) + if err != nil { + t.Fatalf("RunBuild() error = %v", err) + } + if build.Status != BuildStatusSucceeded || build.ArtifactID == "" { + t.Fatalf("build = %+v, want succeeded with artifact", build) + } + artifact, err := manager.Artifact(context.Background(), build.ArtifactID) + if err != nil { + t.Fatalf("Artifact() error = %v", err) + } + if artifact.ArtifactType != ArtifactTypeBinary || artifact.Status != ArtifactStatusLoadable { + t.Fatalf("artifact = %+v, want loadable binary", artifact) + } + if _, err := manager.SetDesired(context.Background(), "admin", artifact.PluginID, artifact.ID, DesiredEnabled, `{}`, 10); err != nil { + t.Fatalf("SetDesired() error = %v", err) + } + plugin, err := manager.Enable(context.Background(), "admin", artifact.PluginID) + if err != nil { + t.Fatalf("Enable() error = %v", err) + } + if plugin.RuntimeState != RuntimeEnabled || plugin.ActiveArtifactID != artifact.ID { + t.Fatalf("plugin = %+v, want enabled built artifact", plugin) + } +} + +func TestManagerBuildFailureDoesNotChangeActiveArtifact(t *testing.T) { + adapter := &fakeAdapter{} + builder := &fakeBuilder{err: errors.New("compile failed")} + manager := newManagerForTestWithBuilders(t, adapter, map[string]SourceBuilder{ + BuilderTypeLocalProcess: builder, + }) + active := uploadTestArtifact(t, manager, "plugin-a") + if _, err := manager.SetDesired(context.Background(), "admin", "plugin-a", active.ID, DesiredEnabled, `{}`, 10); err != nil { + t.Fatalf("SetDesired() error = %v", err) + } + if _, err := manager.Enable(context.Background(), "admin", "plugin-a"); err != nil { + t.Fatalf("Enable() error = %v", err) + } + source := uploadTestSource(t, manager, "plugin-a") + builds, err := manager.ListBuilds(context.Background(), "plugin-a") + if err != nil { + t.Fatalf("ListBuilds() error = %v", err) + } + if len(builds) == 0 || builds[0].SourceID != source.ID { + t.Fatalf("queued builds = %+v, want source %s", builds, source.ID) + } + build, err := manager.RunBuild(context.Background(), "admin", builds[0].ID) + if err != nil { + t.Fatalf("RunBuild() unexpected manager error = %v", err) + } + if build.Status != BuildStatusFailed { + t.Fatalf("build status = %q, want failed", build.Status) + } + plugin, err := manager.Plugin(context.Background(), "plugin-a") + if err != nil { + t.Fatalf("Plugin() error = %v", err) + } + if plugin.ActiveArtifactID != active.ID { + t.Fatalf("active artifact = %q, want unchanged %q", plugin.ActiveArtifactID, active.ID) + } +} + +func TestManagerRejectsSourceArtifactLoad(t *testing.T) { + manager := newManagerForTest(t, &fakeAdapter{}) + source := uploadTestSource(t, manager, "plugin-a") + if _, err := manager.SetDesired(context.Background(), "admin", "plugin-a", source.ID, DesiredEnabled, `{}`, 10); err == nil || !strings.Contains(err.Error(), "binary artifact") { + t.Fatalf("SetDesired(source) error = %v, want binary artifact rejection", err) + } +} + func TestManagerEnableDisableAndDispatch(t *testing.T) { adapter := &fakeAdapter{} manager := newManagerForTest(t, adapter) @@ -455,12 +555,18 @@ func enableProtocolProxyTestPlugin(t *testing.T, manager *Manager, pluginID stri } func newManagerForTest(t *testing.T, adapter RuntimeAdapter) *Manager { + t.Helper() + return newManagerForTestWithBuilders(t, adapter, nil) +} + +func newManagerForTestWithBuilders(t *testing.T, adapter RuntimeAdapter, builders map[string]SourceBuilder) *Manager { t.Helper() db := openPluginManagerTestDB(t) return New(Options{ DB: db, ArtifactRoot: t.TempDir(), Adapter: adapter, + Builders: builders, }) } @@ -499,6 +605,137 @@ func uploadTestArtifactWithCapabilities(t *testing.T, manager *Manager, pluginID return artifact } +func uploadTestSource(t *testing.T, manager *Manager, pluginID string) ArtifactRecord { + t.Helper() + packagePath := writeTestMCGP(t, map[string][]byte{ + "manifest.json": testSourceManifestBytes(t, pluginID), + "go.mod": []byte("module example.com/" + pluginID + "\n\ngo 1.24.0\n"), + "main.go": []byte("package main\n"), + }) + source, err := manager.UploadSource(context.Background(), ArtifactUpload{ + SourcePath: packagePath, + FileName: pluginID + "-source.mcgp", + Actor: "admin", + }) + if err != nil { + t.Fatalf("UploadSource(%s) error = %v", pluginID, err) + } + return source +} + +func uploadBuildableTestSource(t *testing.T, manager *Manager, pluginID string) ArtifactRecord { + t.Helper() + repoRoot, err := filepath.Abs("../..") + if err != nil { + t.Fatalf("Abs(repo root) error = %v", err) + } + manifest := Manifest{ + SchemaVersion: SchemaVersion, + ID: pluginID, + Name: "Buildable Source", + Version: "0.1.0", + ArtifactType: ArtifactTypeSource, + Runtime: RuntimeManifest{ + Type: RuntimeGoPlugin, + EntrySymbol: "Plugin", + MetadataSymbol: "MCGatewayPluginMetadata", + }, + Build: BuildManifest{ + Type: BuildTypeGo, + Entry: ".", + GoVersion: runtime.Version(), + Tags: []string{}, + VendorRequired: false, + Output: RuntimeEntry, + }, + APIVersion: APIVersion, + SDKModule: "github.com/tursom/mc-gateway/plugin/api", + SDKModuleVersion: "v0.1.0", + GoVersion: runtime.Version(), + GOOS: runtime.GOOS, + GOARCH: runtime.GOARCH, + ExtensionPoints: []ExtensionPoint{{ + Type: "hook", + Key: ExtensionUpstreamConnect, + }}, + Capabilities: json.RawMessage(`{"extension_points":["upstream.connect/v1"]}`), + } + manifestBytes, err := json.Marshal(manifest) + if err != nil { + t.Fatalf("Marshal manifest error = %v", err) + } + mainSource := `package main + +import ( + "encoding/json" + "net" + "runtime" + + "github.com/tursom/mc-gateway/plugin/api" +) + +type pluginImpl struct{ api.AbstractPlugin } + +func Plugin() api.Plugin { return &pluginImpl{} } + +func MCGatewayPluginMetadata() string { return manifestJSON } + +func (p *pluginImpl) Init(gateway api.Gateway) error { + return api.RegisterHookHandler( + gateway, + api.HookUpstreamConnect, + func(api.UpstreamConnectRequest) bool { return true }, + func(api.UpstreamConnectRequest) (net.Conn, error) { + left, right := net.Pipe() + _ = right.Close() + return left, nil + }, + ) +} + +var manifestJSON = compactJSON(map[string]any{ + "schema_version": "mc-gateway.plugin/v1", + "id": "` + pluginID + `", + "name": "Buildable Source", + "version": "0.1.0", + "artifact_type": "binary", + "runtime": map[string]any{ + "type": "go-plugin", + "entry": "plugin.so", + "entry_symbol": "Plugin", + "metadata_symbol": "MCGatewayPluginMetadata", + }, + "api_version": "plugin-api/v1", + "sdk_module": "github.com/tursom/mc-gateway/plugin/api", + "sdk_module_version": "v0.1.0", + "go_version": runtime.Version(), + "go_os": runtime.GOOS, + "go_arch": runtime.GOARCH, + "extension_points": []map[string]any{{"type": "hook", "key": "upstream.connect/v1"}}, + "capabilities": map[string]any{"extension_points": []string{"upstream.connect/v1"}}, +}) + +func compactJSON(value any) string { + data, _ := json.Marshal(value) + return string(data) +} +` + packagePath := writeTestMCGP(t, map[string][]byte{ + "manifest.json": manifestBytes, + "go.mod": []byte("module example.com/" + pluginID + "\n\ngo 1.24.0\n\nrequire github.com/tursom/mc-gateway v0.0.0\n\nreplace github.com/tursom/mc-gateway => " + filepath.ToSlash(repoRoot) + "\n"), + "main.go": []byte(mainSource), + }) + source, err := manager.UploadSource(context.Background(), ArtifactUpload{ + SourcePath: packagePath, + FileName: pluginID + "-source.mcgp", + Actor: "admin", + }) + if err != nil { + t.Fatalf("UploadSource(%s) error = %v", pluginID, err) + } + return source +} + func waitForPluginManagerTest(t *testing.T, done func() bool) { t.Helper() deadline := time.After(2 * time.Second) @@ -552,6 +789,15 @@ type fakePlugin struct { api.AbstractPlugin } +type fakeBuilder struct { + result BuildResult + err error +} + +func (b *fakeBuilder) Build(context.Context, ArtifactRecord, BuildRequest, BuildRecord) (BuildResult, error) { + return b.result, b.err +} + func newMemoryConn() net.Conn { left, right := net.Pipe() _ = right.Close() diff --git a/internal/pluginmanager/repository.go b/internal/pluginmanager/repository.go index 62c89bb..a227d76 100644 --- a/internal/pluginmanager/repository.go +++ b/internal/pluginmanager/repository.go @@ -5,6 +5,7 @@ import ( "database/sql" "encoding/json" "errors" + "fmt" "time" ) @@ -120,6 +121,9 @@ func (r Repository) UpsertDesired(ctx context.Context, actor, pluginID, artifact if artifact.PluginID != pluginID { return PluginRecord{}, errors.New("artifact plugin_id does not match") } + if artifact.ArtifactType != ArtifactTypeBinary { + return PluginRecord{}, errors.New("desired artifact must be a binary artifact") + } now := r.now().Unix() tx, err := r.db.BeginTx(ctx, nil) @@ -266,6 +270,207 @@ func (r Repository) UpdateArtifactStatus(ctx context.Context, artifactID, status return err } +func (r Repository) CreateBuild(ctx context.Context, build BuildRecord) (BuildRecord, error) { + now := r.now().Unix() + build.CreatedAt = now + build.UpdatedAt = now + if build.Status == "" { + build.Status = BuildStatusQueued + } + _, err := r.db.ExecContext(ctx, ` +INSERT INTO plugin_builds( + plugin_id, source_id, artifact_id, status, builder_type, builder_image, builder_version, + go_version, go_os, go_arch, go_amd64, go_arm64, cgo_enabled, build_tags, + sdk_module, sdk_version, go_proxy, go_no_sumdb, go_private, vendor_required, + source_sha256, artifact_sha256, module_summary_json, go_version_m_json, abi_fingerprint, + log_summary, metadata_json, error, started_at, ended_at, duration_ms, created_by, created_at, updated_at +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)`, + build.PluginID, build.SourceID, build.ArtifactID, build.Status, build.BuilderType, build.BuilderImage, build.BuilderVersion, + build.GoVersion, build.GOOS, build.GOARCH, build.GOAMD64, build.GOARM64, build.CGOEnabled, build.BuildTags, + build.SDKModule, build.SDKVersion, build.GOPROXY, build.GONOSUMDB, build.GOPRIVATE, boolInt(build.VendorRequired), + build.SourceSHA256, build.ArtifactSHA256, defaultJSONArray(build.ModuleSummary), defaultJSONObject(build.GoVersionM), build.ABIFingerprint, + build.LogSummary, defaultJSONObject(build.MetadataJSON), build.Error, build.StartedAt, build.EndedAt, build.DurationMS, build.CreatedBy, build.CreatedAt, build.UpdatedAt) + if err != nil { + return BuildRecord{}, err + } + id, err := lastInsertID(ctx, r.db) + if err != nil { + return BuildRecord{}, err + } + return r.Build(ctx, id) +} + +func (r Repository) Build(ctx context.Context, id int64) (BuildRecord, error) { + row := r.db.QueryRowContext(ctx, ` +SELECT id, plugin_id, source_id, artifact_id, status, builder_type, builder_image, builder_version, + go_version, go_os, go_arch, go_amd64, go_arm64, cgo_enabled, build_tags, + sdk_module, sdk_version, go_proxy, go_no_sumdb, go_private, vendor_required, + source_sha256, artifact_sha256, module_summary_json, go_version_m_json, abi_fingerprint, + log_summary, metadata_json, error, started_at, ended_at, duration_ms, created_by, created_at, updated_at +FROM plugin_builds WHERE id = ?`, id) + return scanBuild(row) +} + +func (r Repository) ListBuilds(ctx context.Context, pluginID string) ([]BuildRecord, error) { + query := ` +SELECT id, plugin_id, source_id, artifact_id, status, builder_type, builder_image, builder_version, + go_version, go_os, go_arch, go_amd64, go_arm64, cgo_enabled, build_tags, + sdk_module, sdk_version, go_proxy, go_no_sumdb, go_private, vendor_required, + source_sha256, artifact_sha256, module_summary_json, go_version_m_json, abi_fingerprint, + log_summary, metadata_json, error, started_at, ended_at, duration_ms, created_by, created_at, updated_at +FROM plugin_builds` + var args []any + if pluginID != "" { + query += ` WHERE plugin_id = ?` + args = append(args, pluginID) + } + query += ` ORDER BY created_at DESC, id DESC` + rows, err := r.db.QueryContext(ctx, query, args...) + if err != nil { + return nil, err + } + defer rows.Close() + var builds []BuildRecord + for rows.Next() { + build, err := scanBuild(rows) + if err != nil { + return nil, err + } + builds = append(builds, build) + } + return builds, rows.Err() +} + +func (r Repository) MarkBuildRunning(ctx context.Context, id int64) error { + now := r.now().Unix() + _, err := r.db.ExecContext(ctx, `UPDATE plugin_builds SET status = ?, started_at = ?, updated_at = ? WHERE id = ? AND status = ?`, + BuildStatusRunning, now, now, id, BuildStatusQueued) + return err +} + +func (r Repository) FinishBuild(ctx context.Context, build BuildRecord) error { + now := r.now().Unix() + if build.EndedAt == 0 { + build.EndedAt = now + } + if build.StartedAt > 0 && build.DurationMS == 0 { + build.DurationMS = (build.EndedAt - build.StartedAt) * 1000 + } + _, err := r.db.ExecContext(ctx, ` +UPDATE plugin_builds SET + artifact_id = ?, status = ?, builder_image = ?, builder_version = ?, go_version = ?, + go_os = ?, go_arch = ?, go_amd64 = ?, go_arm64 = ?, cgo_enabled = ?, build_tags = ?, + sdk_module = ?, sdk_version = ?, go_proxy = ?, go_no_sumdb = ?, go_private = ?, vendor_required = ?, + source_sha256 = ?, artifact_sha256 = ?, module_summary_json = ?, go_version_m_json = ?, + abi_fingerprint = ?, log_summary = ?, metadata_json = ?, error = ?, ended_at = ?, duration_ms = ?, updated_at = ? +WHERE id = ?`, + build.ArtifactID, build.Status, build.BuilderImage, build.BuilderVersion, build.GoVersion, + build.GOOS, build.GOARCH, build.GOAMD64, build.GOARM64, build.CGOEnabled, build.BuildTags, + build.SDKModule, build.SDKVersion, build.GOPROXY, build.GONOSUMDB, build.GOPRIVATE, boolInt(build.VendorRequired), + build.SourceSHA256, build.ArtifactSHA256, defaultJSONArray(build.ModuleSummary), defaultJSONObject(build.GoVersionM), + build.ABIFingerprint, build.LogSummary, defaultJSONObject(build.MetadataJSON), build.Error, build.EndedAt, build.DurationMS, now, + build.ID) + return err +} + +func (r Repository) CancelBuild(ctx context.Context, id int64, actor string) (BuildRecord, error) { + now := r.now().Unix() + res, err := r.db.ExecContext(ctx, ` +UPDATE plugin_builds +SET status = ?, error = ?, ended_at = ?, updated_at = ? +WHERE id = ? AND status IN (?, ?)`, + BuildStatusCanceled, "build canceled by "+actor, now, now, id, BuildStatusQueued, BuildStatusRunning) + if err != nil { + return BuildRecord{}, err + } + if rows, _ := res.RowsAffected(); rows == 0 { + return r.Build(ctx, id) + } + return r.Build(ctx, id) +} + +func (r Repository) RetryBuild(ctx context.Context, id int64, actor string) (BuildRecord, error) { + previous, err := r.Build(ctx, id) + if err != nil { + return BuildRecord{}, err + } + if previous.Status != BuildStatusFailed && previous.Status != BuildStatusCanceled { + return BuildRecord{}, fmt.Errorf("build status %q cannot be retried", previous.Status) + } + previous.ID = 0 + previous.ArtifactID = "" + previous.ArtifactSHA256 = "" + previous.Status = BuildStatusQueued + previous.LogSummary = "" + previous.Error = "" + previous.StartedAt = 0 + previous.EndedAt = 0 + previous.DurationMS = 0 + previous.CreatedBy = actor + return r.CreateBuild(ctx, previous) +} + +func (r Repository) ReferencedArtifactIDs(ctx context.Context) (map[string]bool, error) { + refs := make(map[string]bool) + rows, err := r.db.QueryContext(ctx, ` +SELECT desired_artifact_id, active_artifact_id, loaded_artifact_id +FROM plugins +WHERE desired_state <> 'deleted'`) + if err != nil { + return nil, err + } + for rows.Next() { + var desired, active, loaded string + if err := rows.Scan(&desired, &active, &loaded); err != nil { + rows.Close() + return nil, err + } + for _, id := range []string{desired, active, loaded} { + if id != "" { + refs[id] = true + } + } + } + if err := rows.Close(); err != nil { + return nil, err + } + rows, err = r.db.QueryContext(ctx, `SELECT artifact_id FROM plugin_config_snapshots WHERE artifact_id <> ''`) + if err != nil { + return nil, err + } + for rows.Next() { + var id string + if err := rows.Scan(&id); err != nil { + rows.Close() + return nil, err + } + if id != "" { + refs[id] = true + } + } + if err := rows.Close(); err != nil { + return nil, err + } + rows, err = r.db.QueryContext(ctx, `SELECT artifact_id FROM plugin_builds WHERE artifact_id <> '' AND status = ?`, BuildStatusSucceeded) + if err != nil { + return nil, err + } + for rows.Next() { + var id string + if err := rows.Scan(&id); err != nil { + rows.Close() + return nil, err + } + if id != "" { + refs[id] = true + } + } + if err := rows.Close(); err != nil { + return nil, err + } + return refs, nil +} + func (r Repository) RecordOperation(ctx context.Context, pluginID, artifactID, operation, status, actor, message string, metadata any) error { metadataJSON, err := marshalDefaultObject(metadata) if err != nil { @@ -322,6 +527,20 @@ func scanPluginRow(row rowScanner, plugin *PluginRecord) error { ) } +func scanBuild(row rowScanner) (BuildRecord, error) { + var build BuildRecord + var vendorRequired int + err := row.Scan( + &build.ID, &build.PluginID, &build.SourceID, &build.ArtifactID, &build.Status, &build.BuilderType, &build.BuilderImage, &build.BuilderVersion, + &build.GoVersion, &build.GOOS, &build.GOARCH, &build.GOAMD64, &build.GOARM64, &build.CGOEnabled, &build.BuildTags, + &build.SDKModule, &build.SDKVersion, &build.GOPROXY, &build.GONOSUMDB, &build.GOPRIVATE, &vendorRequired, + &build.SourceSHA256, &build.ArtifactSHA256, &build.ModuleSummary, &build.GoVersionM, &build.ABIFingerprint, + &build.LogSummary, &build.MetadataJSON, &build.Error, &build.StartedAt, &build.EndedAt, &build.DurationMS, &build.CreatedBy, &build.CreatedAt, &build.UpdatedAt, + ) + build.VendorRequired = vendorRequired != 0 + return build, err +} + func marshalDefaultObject(value any) (string, error) { if value == nil { return "{}", nil @@ -335,3 +554,33 @@ func marshalDefaultObject(value any) (string, error) { } return string(data), nil } + +func defaultJSONObject(value string) string { + if value == "" || !json.Valid([]byte(value)) { + return "{}" + } + return value +} + +func defaultJSONArray(value string) string { + if value == "" || !json.Valid([]byte(value)) { + return "[]" + } + return value +} + +func boolInt(value bool) int { + if value { + return 1 + } + return 0 +} + +func lastInsertID(ctx context.Context, db *sql.DB) (int64, error) { + row := db.QueryRowContext(ctx, `SELECT last_insert_rowid()`) + var id int64 + if err := row.Scan(&id); err != nil { + return 0, err + } + return id, nil +} diff --git a/internal/pluginmanager/types.go b/internal/pluginmanager/types.go index 83a5a2c..b3a4153 100644 --- a/internal/pluginmanager/types.go +++ b/internal/pluginmanager/types.go @@ -15,8 +15,10 @@ const ( APIVersion = "plugin-api/v1" ArtifactTypeBinary = "binary" + ArtifactTypeSource = "source" RuntimeGoPlugin = "go-plugin" RuntimeEntry = "plugin.so" + SourceBuildEntry = "." ExtensionUpstreamConnect = "upstream.connect/v1" @@ -30,6 +32,16 @@ const ( ArtifactStatusRejected = "rejected" ArtifactStatusDeleted = "deleted" + BuildStatusQueued = "queued" + BuildStatusRunning = "running" + BuildStatusSucceeded = "succeeded" + BuildStatusFailed = "failed" + BuildStatusCanceled = "canceled" + + BuilderTypeLocalProcess = "local-process" + BuilderTypeContainer = "container" + BuildTypeGo = "go" + DesiredEnabled = "enabled" DesiredDisabled = "disabled" DesiredDeleted = "deleted" @@ -49,6 +61,7 @@ const ( DefaultExtractedMaxBytes = 256 * 1024 * 1024 DefaultNonRuntimeMaxBytes = 16 * 1024 * 1024 DefaultInitialWriteTimeout = time.Second + DefaultBuildLogMaxBytes = 64 * 1024 ) var ( @@ -64,6 +77,7 @@ type Manifest struct { Description string `json:"description"` ArtifactType string `json:"artifact_type"` Runtime RuntimeManifest `json:"runtime"` + Build BuildManifest `json:"build,omitempty"` APIVersion string `json:"api_version"` SDKModule string `json:"sdk_module"` SDKModuleVersion string `json:"sdk_module_version"` @@ -80,10 +94,21 @@ type Manifest struct { type RuntimeManifest struct { Type string `json:"type"` Entry string `json:"entry"` + BuildEntry string `json:"build_entry"` EntrySymbol string `json:"entry_symbol"` MetadataSymbol string `json:"metadata_symbol"` } +type BuildManifest struct { + Type string `json:"type"` + Entry string `json:"entry"` + GoVersion string `json:"go_version"` + CGOEnabled *bool `json:"cgo_enabled,omitempty"` + Tags []string `json:"tags"` + VendorRequired bool `json:"vendor_required"` + Output string `json:"output"` +} + type ExtensionPoint struct { Type string `json:"type"` Key string `json:"key"` @@ -183,6 +208,75 @@ type OperationRecord struct { CreatedAt int64 `json:"created_at"` } +type BuildRecord struct { + ID int64 `json:"id"` + PluginID string `json:"plugin_id"` + SourceID string `json:"source_id"` + ArtifactID string `json:"artifact_id"` + Status string `json:"status"` + BuilderType string `json:"builder_type"` + BuilderImage string `json:"builder_image"` + BuilderVersion string `json:"builder_version"` + GoVersion string `json:"go_version"` + GOOS string `json:"go_os"` + GOARCH string `json:"go_arch"` + GOAMD64 string `json:"go_amd64"` + GOARM64 string `json:"go_arm64"` + CGOEnabled string `json:"cgo_enabled"` + BuildTags string `json:"build_tags"` + SDKModule string `json:"sdk_module"` + SDKVersion string `json:"sdk_version"` + GOPROXY string `json:"go_proxy"` + GONOSUMDB string `json:"go_no_sumdb"` + GOPRIVATE string `json:"go_private"` + VendorRequired bool `json:"vendor_required"` + SourceSHA256 string `json:"source_sha256"` + ArtifactSHA256 string `json:"artifact_sha256"` + ModuleSummary string `json:"module_summary_json"` + GoVersionM string `json:"go_version_m_json"` + ABIFingerprint string `json:"abi_fingerprint"` + LogSummary string `json:"log_summary"` + MetadataJSON string `json:"metadata_json"` + Error string `json:"error"` + StartedAt int64 `json:"started_at"` + EndedAt int64 `json:"ended_at"` + DurationMS int64 `json:"duration_ms"` + CreatedBy string `json:"created_by"` + CreatedAt int64 `json:"created_at"` + UpdatedAt int64 `json:"updated_at"` +} + +type BuildRequest struct { + SourceID string `json:"source_id"` + BuilderType string `json:"builder_type"` + BuilderImage string `json:"builder_image"` + BuilderVersion string `json:"builder_version"` + GOOS string `json:"go_os"` + GOARCH string `json:"go_arch"` + GOAMD64 string `json:"go_amd64"` + GOARM64 string `json:"go_arm64"` + CGOEnabled string `json:"cgo_enabled"` + BuildTags string `json:"build_tags"` + SDKModule string `json:"sdk_module"` + SDKVersion string `json:"sdk_version"` + GOPROXY string `json:"go_proxy"` + GONOSUMDB string `json:"go_no_sumdb"` + GOPRIVATE string `json:"go_private"` + VendorRequired bool `json:"vendor_required"` +} + +type GCCandidate struct { + Kind string `json:"kind"` + ID string `json:"id"` + PluginID string `json:"plugin_id"` + Path string `json:"path"` + Protected bool `json:"protected"` + Reason string `json:"reason"` + SizeBytes int64 `json:"size_bytes"` + CreatedAt int64 `json:"created_at"` + Referenced bool `json:"referenced"` +} + type ConfigSnapshot struct { ID int64 `json:"id"` PluginID string `json:"plugin_id"`