diff --git a/.github/workflows/autotest_prs.yml b/.github/workflows/autotest_prs.yml index d87ddc05..d64c17ca 100644 --- a/.github/workflows/autotest_prs.yml +++ b/.github/workflows/autotest_prs.yml @@ -35,7 +35,7 @@ jobs: staticcheck ./... - name: Set up MinIO - uses: infleet/minio-action@v0.0.1 + uses: cohere-llc/minio-action@v0.0.3 with: port: "9000" version: "latest" diff --git a/.github/workflows/irods.yml b/.github/workflows/irods.yml index 2d65f0ea..3603e66d 100644 --- a/.github/workflows/irods.yml +++ b/.github/workflows/irods.yml @@ -26,7 +26,7 @@ jobs: uses: actions/checkout@v4 - name: Set up MinIO - uses: infleet/minio-action@v0.0.1 + uses: cohere-llc/minio-action@v0.0.3 with: port: "9000" version: "latest" diff --git a/README.md b/README.md index 447ba68d..113bc3e0 100644 --- a/README.md +++ b/README.md @@ -45,7 +45,7 @@ require a Minio test instance to be running. You can start one with docker or podman: ``` -docker run -d -p 9000:9000 -p 9001:9001 -e "MINIO_ROOT_USER=minioadmin" -e "MINIO_ROOT_PASSWORD=minioadmin" minio/minio server /data --console-address ":9001" +docker run -d -p 9000:9000 -p 9001:9001 -e "MINIO_ROOT_USER=minioadmin" -e "MINIO_ROOT_PASSWORD=minioadmin" pgsty/minio server /data --console-address ":9001" ``` Then you can run these tests as you would any other Go project: @@ -87,4 +87,4 @@ to do: authenticate with the JGI Data Portal Alternatively, you can run tests against mock services without the above -environment variables set by setting `DTS_TEST_WITH_MOCK_SERVICES=true` \ No newline at end of file +environment variables set by setting `DTS_TEST_WITH_MOCK_SERVICES=true` diff --git a/auth/auth.go b/auth/auth.go index e1dda5a0..6a079863 100644 --- a/auth/auth.go +++ b/auth/auth.go @@ -33,10 +33,14 @@ type User struct { Organization string // true if this user is a Superuser IsSuper bool + // credentials for connections between endpoints with different providers (e.g. Globus <--> S3) + ConnectionCredentials map[string]Credential } // A credential used for authorization and authentication type Credential struct { + // the username associated with this credential + Username string `yaml:"username"` // the ID used for authorization (username or UUID) Id string `yaml:"id"` // the secret used for authentication (e.g. password) diff --git a/auth/authenticator.go b/auth/authenticator.go index fef116a2..af48d27b 100644 --- a/auth/authenticator.go +++ b/auth/authenticator.go @@ -143,11 +143,12 @@ func (a *Authenticator) readAccessTokenFile() error { } userRecords[token] = User{ - Name: record[0], - Email: record[1], - Orcid: record[2], - Organization: record[3], - IsSuper: isSuper, + Name: record[0], + Email: record[1], + Orcid: record[2], + Organization: record[3], + IsSuper: isSuper, + ConnectionCredentials: make(map[string]Credential), } } diff --git a/auth/kbase_auth_server.go b/auth/kbase_auth_server.go index 561405cb..fae40353 100644 --- a/auth/kbase_auth_server.go +++ b/auth/kbase_auth_server.go @@ -72,10 +72,10 @@ func NewKBaseAuthServer(accessToken string, options ...KBaseAuthServerOption) (* } // check our list of KBase auth server instances for this access token - if instances == nil { - instances = make(map[string]*KBaseAuthServer) + if instances_ == nil { + instances_ = make(map[string]*KBaseAuthServer) } - if server, found := instances[accessToken]; found { + if server, found := instances_[accessToken]; found { return server, nil } else { server := KBaseAuthServer{ @@ -91,7 +91,7 @@ func NewKBaseAuthServer(accessToken string, options ...KBaseAuthServerOption) (* } // register this instance of the auth server - instances[accessToken] = &server + instances_[accessToken] = &server return &server, err } } @@ -103,8 +103,9 @@ func (server KBaseAuthServer) User() (User, error) { return User{}, err } user := User{ - Name: kbUser.Display, - Email: kbUser.Email, + Name: kbUser.Display, + Email: kbUser.Email, + ConnectionCredentials: make(map[string]Credential), } for _, pid := range kbUser.Idents { // grab the first ORCID associated with the user @@ -113,6 +114,23 @@ func (server KBaseAuthServer) User() (User, error) { break } } + + // try to access the MMS in case we're talking to the KBase Lakehouse + mms := NewMMS() + record, err := mms.FetchRecord(server.AccessToken) + if err == nil { + user.ConnectionCredentials["s3"] = Credential{ + Username: record.Username, + Id: record.S3AccessKey, + Secret: record.S3SecretKey, + } + user.ConnectionCredentials["polaris"] = Credential{ + Username: record.Username, + Id: record.PolarisClientId, + Secret: record.PolarisClientSecret, + } + } + return user, nil } @@ -155,7 +173,7 @@ type kbaseAuthErrorResponse struct { // here's a set of instances to the KBase auth server, mapped by OAuth2 // access token -var instances map[string]*KBaseAuthServer +var instances_ map[string]*KBaseAuthServer // emits an error representing the error in a response to the auth server func kbaseAuthError(response *http.Response) error { diff --git a/auth/kbase_mms.go b/auth/kbase_mms.go new file mode 100644 index 00000000..3a6c5c03 --- /dev/null +++ b/auth/kbase_mms.go @@ -0,0 +1,85 @@ +// Copyright (c) 2023 The KBase Project and its Contributors +// Copyright (c) 2023 Cohere Consulting, LLC +// +// Permission is hereby granted, free of charge, to any person obtaining a copy of +// this software and associated documentation files (the "Software"), to deal in +// the Software without restriction, including without limitation the rights to +// use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies +// of the Software, and to permit persons to whom the Software is furnished to do +// so, subject to the following conditions: +// +// The above copyright notice and this permission notice shall be included in all +// copies or substantial portions of the Software. +// +// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +// SOFTWARE. + +package auth + +import ( + "encoding/json" + "fmt" + "io" + "net/http" + "time" +) + +// The Minio Management Service (MMS) provide authentication information for a user +// given a valid KBase token for that user + +type MMSRecord struct { + Username string `json:"username"` + S3AccessKey string `json:"s3_access_key"` + S3SecretKey string `json:"s3_secret_key"` + PolarisClientId string `json:"polaris_client_id"` + PolarisClientSecret string `json:"polaris_client_secret"` +} + +type MMS struct { + Client http.Client +} + +func NewMMS() MMS { + return MMS{ + Client: http.Client{ + Timeout: 5 * time.Second, + }, + } +} + +// retrieves the MMS record associated with the given access token +func (mms MMS) FetchRecord(accessToken string) (MMSRecord, error) { + resource := fmt.Sprintf("%s:%d", kbaseMMSUrl, kbaseMMSPort) + "/credentials/" + request, err := http.NewRequest(http.MethodGet, resource, http.NoBody) + if err != nil { + return MMSRecord{}, err + } + request.Header.Add("Authorization", fmt.Sprintf("Bearer %s", accessToken)) + resp, err := mms.Client.Do(request) + if err != nil { + return MMSRecord{}, err + } + defer resp.Body.Close() + if resp.StatusCode < http.StatusOK || resp.StatusCode >= http.StatusMultipleChoices { + return MMSRecord{}, fmt.Errorf("MMS returned HTTP %d", resp.StatusCode) + } + + body, err := io.ReadAll(resp.Body) + if err != nil { + return MMSRecord{}, err + } + + var record MMSRecord + err = json.Unmarshal(body, &record) + return record, err +} + +const ( + kbaseMMSUrl = "http://mms.dev" + kbaseMMSPort = 8000 +) diff --git a/databases/kbase/database.go b/databases/kbase/database.go index 9efa4283..c9f0303b 100644 --- a/databases/kbase/database.go +++ b/databases/kbase/database.go @@ -52,7 +52,7 @@ func NewDatabase(conf Config) (databases.Database, error) { EndpointName: conf.Endpoint, } var err error - db.kbaseFed, err = newKBaseUserFederation(conf.KBaseUserFederationConfig) + db.kbaseFed, err = NewKBaseUserFederation(conf.KBaseUserFederationConfig) if err != nil { return nil, err } @@ -107,7 +107,7 @@ func (db *Database) Finalize(orcid string, id uuid.UUID) error { } func (db *Database) LocalUser(orcid string) (string, error) { - return db.kbaseFed.usernameForOrcid(orcid) + return db.kbaseFed.UsernameForOrcid(orcid) } func (db Database) Save() (databases.DatabaseSaveState, error) { diff --git a/databases/kbase/user_federation.go b/databases/kbase/user_federation.go index 5c4b4faa..76b6661b 100644 --- a/databases/kbase/user_federation.go +++ b/databases/kbase/user_federation.go @@ -31,27 +31,29 @@ import ( "strings" "time" "unicode" + + "github.com/google/uuid" ) //======================= // KBase user federation //======================= -// In order to map an ORCID to a KBase username, we maintain a mapping that -// stores entries for all KBase users with ORCIDs. This mapping currently lives -// a 2-column spreadsheet (CSV) in the DTS data directory. The data in this -// spreadsheet is reloaded every hour on the top of the hour so a new file can -// be dropped into the data directory with predictable results. +// In order to map an ORCID to a KBase username and Globus ID, we maintain a mapping that stores +// entries for all KBase users with associated ORCIDs and Globus IDs. This mapping currently lives +// in a 3-column spreadsheet (CSV) in the DTS data directory. Column names are The data in this spreadsheet is +// reloaded every hour on the top of the hour so a new file can be dropped into the data directory +// with predictable results. Extra columns are ignored. // KBase User Federation type KBaseUserFederation struct { Started bool - FilePath string // full path to the KBase user table file - UpdateChan chan struct{} // triggers updates to the ORCID/user table - StopChan chan struct{} // stops the user federation subsystem - OrcidChan chan string // passes ORCIDs in for lookup - UserChan chan string // passes usernames out - ErrorChan chan error // passes errors out + FilePath string // full path to the KBase user table file + UpdateChan chan struct{} // triggers updates to the ORCID/user table + StopChan chan struct{} // stops the user federation subsystem + OrcidChan chan string // passes ORCIDs in for lookup + RecordChan chan kbaseUserRecord // passes user info out + ErrorChan chan error // passes errors out } // configuration information for the KBase user federation subsystem @@ -59,13 +61,20 @@ type KBaseUserFederationConfig struct { DataDirectory string `yaml:"data_directory" mapstructure:"data_directory"` } -func newKBaseUserFederation(conf KBaseUserFederationConfig) (KBaseUserFederation, error) { +func NewKBaseUserFederation(conf KBaseUserFederationConfig) (KBaseUserFederation, error) { kbaseFed := KBaseUserFederation{} kbaseFed.Started = false kbaseFed.FilePath = filepath.Join(conf.DataDirectory, kbaseUserTableFile) return kbaseFed, nil } +func NewKBaseUserFederationFromFile(filename string) KBaseUserFederation { + kbaseFed := KBaseUserFederation{} + kbaseFed.Started = false + kbaseFed.FilePath = filename + return kbaseFed +} + // starts up the user federation machinery if it hasn't yet been started func (kbaseFed *KBaseUserFederation) Start() error { if kbaseFed.Started { @@ -102,14 +111,25 @@ func (kbaseFed *KBaseUserFederation) Start() error { } // returns the KBase username associated with the given ORCID -func (kbaseFed *KBaseUserFederation) usernameForOrcid(orcid string) (string, error) { +func (kbaseFed *KBaseUserFederation) UsernameForOrcid(orcid string) (string, error) { if !kbaseFed.Started { return "", fmt.Errorf("KBase federated user table not available") } kbaseFed.OrcidChan <- orcid - username := <-kbaseFed.UserChan + record := <-kbaseFed.RecordChan err := <-kbaseFed.ErrorChan - return username, err + return record.Username, err +} + +// returns the Globus ID associated with the given ORCID +func (kbaseFed *KBaseUserFederation) GlobusIdForOrcid(orcid string) (uuid.UUID, error) { + if !kbaseFed.Started { + return uuid.UUID{}, fmt.Errorf("KBase federated user table not available") + } + kbaseFed.OrcidChan <- orcid + record := <-kbaseFed.RecordChan + err := <-kbaseFed.ErrorChan + return record.GlobusId, err } func (kbaseFed *KBaseUserFederation) reloadUserTable() error { @@ -133,19 +153,24 @@ func (kbaseFed *KBaseUserFederation) Stop() error { const kbaseUserTableFile = "kbase_user_orcids.csv" +type kbaseUserRecord struct { + Username string + GlobusId uuid.UUID +} + // This goroutine maintains a table that associates ORCIDs with KBase users. // It fields requests for usernames given ORCIDs, and can also update the table // by reading a file. func (kbaseFed *KBaseUserFederation) kbaseUserFederation(started chan struct{}) { // channels kbaseFed.OrcidChan = make(chan string) - kbaseFed.UserChan = make(chan string) + kbaseFed.RecordChan = make(chan kbaseUserRecord) kbaseFed.ErrorChan = make(chan error) kbaseFed.UpdateChan = make(chan struct{}) kbaseFed.StopChan = make(chan struct{}) // mapping of ORCIDs to KBase users - kbaseUserTable := make(map[string]string) + kbaseUserTable := make(map[string]kbaseUserRecord) // we're ready kbaseFed.Started = true @@ -154,11 +179,11 @@ func (kbaseFed *KBaseUserFederation) kbaseUserFederation(started chan struct{}) for { select { case orcid := <-kbaseFed.OrcidChan: // fetching username for orcid - if username, found := kbaseUserTable[orcid]; found { - kbaseFed.UserChan <- username + if record, found := kbaseUserTable[orcid]; found { + kbaseFed.RecordChan <- record kbaseFed.ErrorChan <- nil } else { - kbaseFed.UserChan <- "" + kbaseFed.RecordChan <- kbaseUserRecord{} kbaseFed.ErrorChan <- fmt.Errorf("KBase user not found for ORCID %s", orcid) } case <-kbaseFed.UpdateChan: // update ORCID/user table @@ -180,9 +205,9 @@ type UserOrcidRecord struct { User, Orcid string } -// reads the user table file within the DTS data directory, returning a map -// with ORCID keys associated with username values -func (kbaseFed *KBaseUserFederation) readUserTable() (map[string]string, error) { +// reads the user table file within the DTS data directory, returning a map of ORCID keys to +// KBase user records (with username, perhaps aGlobus ID) +func (kbaseFed *KBaseUserFederation) readUserTable() (map[string]kbaseUserRecord, error) { // open the CVS file containing the user mapping filename := kbaseFed.FilePath slog.Info(fmt.Sprintf("Reading KBase user table from %s", filename)) @@ -199,7 +224,7 @@ func (kbaseFed *KBaseUserFederation) readUserTable() (map[string]string, error) // by a comma. The first line is almost certainly a header with column names, // but we can't be sure, so we simply read every line, checking that // - // * there are 2 entries separated by exactly one comma + // * there are 3 entries separated by exactly one comma // * exactly one of the entries is a well-formed ORCID (xxxx-xxxx-xxxx-xxxx) // * the other entry is a non-empty string with no special characters // @@ -209,16 +234,18 @@ func (kbaseFed *KBaseUserFederation) readUserTable() (map[string]string, error) // requirements is ignored. If there's at least one valid line, we clear the existing KBase user // table and add each (ORCID, user) pair to the user table. - // Finally, there must be a 1:1 correspondence between KBase users and ORCIDs. Otherwise we can't - // map between these items. We track (user, orcid) pairs that violate this constraint and report - // them after we read the entire table. + // Finally, there must be a 1:1 correspondence between KBase users, ORCIDs, and Globus IDs. + // Otherwise we can't map between these items. We track (user, orcid, globus_id) triples that + // violate this constraint and report them after we read the entire table. multipleUsersForOrcid := make(map[string][]string) multipleOrcidsForUser := make(map[string][]string) + multipleGlobusIdsForOrcid := make(map[string][]string) orcidColumn := -1 userColumn := -1 + globusIdColumn := -1 orcidsForUsers := make(map[string]string) - usersForOrcids := make(map[string]string) + recordsForOrcids := make(map[string]kbaseUserRecord) reader := csv.NewReader(file) reader.Comment = '#' records, err := reader.ReadAll() @@ -228,65 +255,102 @@ func (kbaseFed *KBaseUserFederation) readUserTable() (map[string]string, error) Message: "Couldn't parse CVS file", } } - for _, record := range records { - if len(record) != 2 { + for row, record := range records { + if len(record) < 2 { return nil, &InvalidKBaseUserSpreadsheetError{ File: kbaseUserTableFile, - Message: fmt.Sprintf("%d comma-separated columns found (2 expected)", len(record)), + Message: fmt.Sprintf("%d comma-separated columns found (2+ expected)", len(record)), } } - if orcidColumn == -1 { // find the column with an ORCID - for i := range 2 { + // figure out the relevant columns + if orcidColumn == -1 { + for i := range record { if isOrcid(record[i]) { orcidColumn = i - userColumn = (i + 1) % 2 // user column's the other one + } else if len(record) >= 3 && isGlobusId(record[i]) { + globusIdColumn = i + } else if isUsername(record[i]) { + userColumn = i + } + } + if orcidColumn == -1 { + if row == 0 { // first record, ignore + userColumn = -1 + continue + } + return nil, &InvalidKBaseUserSpreadsheetError{ + File: kbaseUserTableFile, + Message: "no ORCID column found", + } + } + if userColumn == -1 { + return nil, &InvalidKBaseUserSpreadsheetError{ + File: kbaseUserTableFile, + Message: "no username column found", + } + } + } + + // keep checking for a Globus ID column if we haven't found it yet + if len(record) >= 3 && globusIdColumn == -1 { + for i := range record { + if isGlobusId(record[i]) { // can't be confused with ORCID or username + globusIdColumn = i } } - } else if !isOrcid(record[orcidColumn]) { - // we've already established the ORCID column, but this line disagrees, - // so the whole file is suspect + } + + if !isOrcid(record[orcidColumn]) || (globusIdColumn != -1 && record[globusIdColumn] != "" && !isGlobusId(record[globusIdColumn])) || !isUsername(record[userColumn]) { + // we've already established the layout, but this line disagrees, so the whole file is suspect return nil, &InvalidKBaseUserSpreadsheetError{ File: kbaseUserTableFile, - Message: "Different lines list username, ORCID data in different columns", + Message: fmt.Sprintf("row %d: Different lines list username, ORCID, globus ID data in different columns", row+1), } } - if orcidColumn != -1 { - orcid := record[orcidColumn] - // ORCID column's okay, but what about the user column? - if !isUsername(record[userColumn]) { - continue + orcid := record[orcidColumn] + username := record[userColumn] + var globusId uuid.UUID + if globusIdColumn != -1 && record[globusIdColumn] != "" { + globusId = uuid.MustParse(record[globusIdColumn]) + } + + // have we seen this ORCID or username/Globus ID before? It's okay, as long as everything + // is consistent + if existingRecord, found := recordsForOrcids[orcid]; found { + if existingRecord.Username != username { + _, found := multipleUsersForOrcid[orcid] + if !found { + multipleUsersForOrcid[orcid] = []string{existingRecord.Username, username} + } } - username := record[userColumn] - - // have we seen this ORCID or username before? It's okay, as long as everything - // is consistent - if existingUser, found := usersForOrcids[orcid]; found { - if existingUser != username { - _, found := multipleUsersForOrcid[orcid] - if !found { - multipleUsersForOrcid[orcid] = []string{existingUser, username} - } + if existingRecord.GlobusId != globusId { + _, found := multipleGlobusIdsForOrcid[orcid] + if !found { + multipleGlobusIdsForOrcid[orcid] = []string{existingRecord.GlobusId.String(), globusId.String()} } - } else { - usersForOrcids[orcid] = username } - if existingOrcid, found := orcidsForUsers[username]; found { - if existingOrcid != orcid { - _, found := multipleOrcidsForUser[username] - if !found { - multipleOrcidsForUser[username] = []string{existingOrcid, orcid} - } + } else { + recordsForOrcids[orcid] = kbaseUserRecord{ + Username: username, + GlobusId: globusId, + } + } + if existingOrcid, found := orcidsForUsers[username]; found { + if existingOrcid != orcid { + _, found := multipleOrcidsForUser[username] + if !found { + multipleOrcidsForUser[username] = []string{existingOrcid, orcid} } - } else { - orcidsForUsers[username] = orcid } + } else { + orcidsForUsers[username] = orcid } } // report any violations of the 1:1 user <-> orcid correspondence - if len(multipleUsersForOrcid) > 0 || len(multipleOrcidsForUser) > 0 { + if len(multipleUsersForOrcid) > 0 || len(multipleOrcidsForUser) > 0 || len(multipleGlobusIdsForOrcid) > 0 { var b strings.Builder for orcid, users := range multipleUsersForOrcid { fmt.Fprintf(&b, "ORCID %s is associated with multiple KBase users: %s\n", orcid, strings.Join(users, ", ")) @@ -294,20 +358,23 @@ func (kbaseFed *KBaseUserFederation) readUserTable() (map[string]string, error) for user, orcids := range multipleOrcidsForUser { fmt.Fprintf(&b, "KBase user %s is associated with multiple ORCIDS: %s\n", user, strings.Join(orcids, ", ")) } + for orcid, globusId := range multipleGlobusIdsForOrcid { + fmt.Fprintf(&b, "ORCID %s is associated with multiple Globus IDs: %s\n", orcid, strings.Join(globusId, ", ")) + } return nil, &InvalidKBaseUserSpreadsheetError{ File: kbaseUserTableFile, - Message: fmt.Sprintf("No 1:1 correspondence exists between users and ORCIDS:\n %s", b.String()), + Message: fmt.Sprintf("No 1:1 correspondence exists between users, ORCIDS, Globus IDs:\n %s", b.String()), } } - if len(usersForOrcids) == 0 { + if len(recordsForOrcids) == 0 { return nil, &InvalidKBaseUserSpreadsheetError{ File: kbaseUserTableFile, Message: "No valid username/ORCID pairs found", } } - return usersForOrcids, nil + return recordsForOrcids, nil } // returns true iff s contains a valid username @@ -322,3 +389,9 @@ func isOrcid(s string) bool { matched, err := regexp.MatchString(`^(\d{4}-){3}\d{3}[\dX]$`, s) return err == nil && matched } + +// returns true iff s contains a valid UUID (nnnnnnnn-nnnn-nnnn-nnnn-nnnnnnnnnnnn) +func isGlobusId(s string) bool { + _, err := uuid.Parse(s) + return err == nil +} diff --git a/databases/kbase/user_federation_test.go b/databases/kbase/user_federation_test.go index cff85d97..91df1c2c 100644 --- a/databases/kbase/user_federation_test.go +++ b/databases/kbase/user_federation_test.go @@ -12,31 +12,18 @@ import ( // valid user table csv contents var goodUserTables = []string{ - `username,orcid -Alice,1234-5678-9101-112X -Bob,1234-5678-9101-1121 -Dave,9402-1876-5432-1098 + `username,orcid,globusid +Alice,1234-5678-9101-112X,184014ac-97c0-4270-94af-97f0cc055673 +Bob,1234-5678-9101-1121,4c5ad8e6-0f5d-4c06-a05d-c6198635675c +Dave,9402-1876-5432-1098, `, - `orcid,username -1234-5678-9101-112X,Alice -1234-5678-9101-1121,Bob -4321-1876-5432-1098,Charlie + `orcid,globusid,username +1234-5678-9101-112X,184014ac-97c0-4270-94af-97f0cc055673,Alice +1234-5678-9101-1121,4c5ad8e6-0f5d-4c06-a05d-c6198635675c,Bob +4321-1876-5432-1098,95f69174-be49-479d-9c71-812580b1371e,Charlie `, } -var goodUserMap = [2]map[string]string{ - { - "1234-5678-9101-112X": "Alice", - "1234-5678-9101-1121": "Bob", - "9402-1876-5432-1098": "Dave", - }, - { - "1234-5678-9101-112X": "Alice", - "1234-5678-9101-1121": "Bob", - "4321-1876-5432-1098": "Charlie", - }, -} - // invalid user table csv contents var badUserTables = []string{ `nocommas`, @@ -116,38 +103,28 @@ func copyDataFile(testDir, src, dst string) error { return err } -func newTestKbaseUserFederation(t *testing.T, filePath string) KBaseUserFederation { - kbaseFed := KBaseUserFederation{ - Started: false, - FilePath: filePath, - UpdateChan: make(chan struct{}), - StopChan: make(chan struct{}), - OrcidChan: make(chan string), - UserChan: make(chan string), - ErrorChan: make(chan error), - } - return kbaseFed -} - func TestKBaseStartReloadStop(t *testing.T) { assert := assert.New(t) - kbaseFed := newTestKbaseUserFederation(t, filepath.Join(testDataDir, "good_user_table_0.csv")) + kbaseFed := NewKBaseUserFederationFromFile(filepath.Join(testDataDir, "good_user_table_0.csv")) err := kbaseFed.Start() assert.Nil(err, "Error starting KBase user federation") - // look up a user - username, err := kbaseFed.usernameForOrcid("1234-5678-9101-112X") + // look up a user and that user's Globus ID + username, err := kbaseFed.UsernameForOrcid("1234-5678-9101-112X") assert.Nil(err, "Error looking up existing ORCID") assert.Equal("Alice", username, "Incorrect username for existing ORCID") + globusId, err := kbaseFed.GlobusIdForOrcid("1234-5678-9101-112X") + assert.Nil(err, "Error looking up existing Globus ID") + assert.Equal("184014ac-97c0-4270-94af-97f0cc055673", globusId.String()) // look up another user - username, err = kbaseFed.usernameForOrcid("9402-1876-5432-1098") + username, err = kbaseFed.UsernameForOrcid("9402-1876-5432-1098") assert.Nil(err, "Error looking up existing ORCID") assert.Equal("Dave", username, "Incorrect username for existing ORCID") // look up a non-existing user - username, err = kbaseFed.usernameForOrcid("9999-8888-7777-6666") + username, err = kbaseFed.UsernameForOrcid("9999-8888-7777-6666") assert.NotNil(err, "No error looking up non-existing ORCID") assert.Equal("", username, "Username returned for non-existing ORCID") @@ -161,17 +138,17 @@ func TestKBaseStartReloadStop(t *testing.T) { assert.Nil(err, "Error reloading user table") // look up a user from the updated table - username, err = kbaseFed.usernameForOrcid("1234-5678-9101-1121") + username, err = kbaseFed.UsernameForOrcid("1234-5678-9101-1121") assert.Nil(err, "Error looking up existing ORCID after reload") assert.Equal("Bob", username, "Incorrect username for existing ORCID after reload") // look up another user from the updated table - username, err = kbaseFed.usernameForOrcid("4321-1876-5432-1098") + username, err = kbaseFed.UsernameForOrcid("4321-1876-5432-1098") assert.Nil(err, "Error looking up existing ORCID after reload") assert.Equal("Charlie", username, "Incorrect username for existing ORCID after reload") // look up an ORCID that existed in the old table but not in the new table - username, err = kbaseFed.usernameForOrcid("9402-1876-5432-1098") + username, err = kbaseFed.UsernameForOrcid("9402-1876-5432-1098") assert.NotNil(err, "No error looking up old ORCID after reload") assert.Equal("", username, "Username returned for old ORCID after reload") @@ -184,111 +161,11 @@ func TestKBaseStartReloadStop(t *testing.T) { assert.NotNil(err, "No error stopping KBase user federation again") // try to look up a user after stopping - username, err = kbaseFed.usernameForOrcid("1234-5678-9101-112X") + username, err = kbaseFed.UsernameForOrcid("1234-5678-9101-112X") assert.NotNil(err, "No error looking up ORCID after stopping federation") assert.Equal("", username, "Username returned after stopping federation") } -func TestKbaseUserFederation(t *testing.T) { - assert := assert.New(t) - - kbaseFed := newTestKbaseUserFederation(t, filepath.Join(testDataDir, "good_user_table_0.csv")) - started := make(chan struct{}) - go kbaseFed.kbaseUserFederation(started) - <-started - - // load the user table - kbaseFed.UpdateChan <- struct{}{} - err := <-kbaseFed.ErrorChan - assert.Nil(err, "Error loading user table") - - // test existing ORCID - kbaseFed.OrcidChan <- "1234-5678-9101-112X" - username := <-kbaseFed.UserChan - err = <-kbaseFed.ErrorChan - assert.Nil(err, "Error looking up existing ORCID") - assert.Equal("Alice", username, "Incorrect username for existing ORCID") - - // test another existing ORCID - kbaseFed.OrcidChan <- "9402-1876-5432-1098" - username = <-kbaseFed.UserChan - err = <-kbaseFed.ErrorChan - assert.Nil(err, "Error looking up existing ORCID") - assert.Equal("Dave", username, "Incorrect username for existing ORCID") - - // test non-existing ORCID - kbaseFed.OrcidChan <- "9999-8888-7777-6666" - username = <-kbaseFed.UserChan - err = <-kbaseFed.ErrorChan - assert.NotNil(err, "No error looking up non-existing ORCID") - assert.Equal("", username, "Username returned for non-existing ORCID") - - // reload user table with updated data - kbaseFed.FilePath = filepath.Join(testDataDir, "good_user_table_1.csv") - kbaseFed.UpdateChan <- struct{}{} - err = <-kbaseFed.ErrorChan - assert.Nil(err, "Error updating user table") - - // test existing ORCID from updated table - kbaseFed.OrcidChan <- "1234-5678-9101-1121" - username = <-kbaseFed.UserChan - err = <-kbaseFed.ErrorChan - assert.Nil(err, "Error looking up existing ORCID after update") - assert.Equal("Bob", username, "Incorrect username for existing ORCID after update") - - // test another existing ORCID from updated table - kbaseFed.OrcidChan <- "4321-1876-5432-1098" - username = <-kbaseFed.UserChan - err = <-kbaseFed.ErrorChan - assert.Nil(err, "Error looking up existing ORCID after update") - assert.Equal("Charlie", username, "Incorrect username for existing ORCID after update") - - // test ORCID that existed in old table but not in new table - kbaseFed.OrcidChan <- "9402-1876-5432-1098" - username = <-kbaseFed.UserChan - err = <-kbaseFed.ErrorChan - assert.NotNil(err, "No error looking up old ORCID after update") - assert.Equal("", username, "Username returned for old ORCID after update") - - // stop the user federation goroutine - kbaseFed.StopChan <- struct{}{} -} - -func TestReadUserTable(t *testing.T) { - assert := assert.New(t) - - for i := range goodUserTables { - filePath := filepath.Join(testDataDir, fmt.Sprintf("good_user_table_%d.csv", i)) - kbaseFed := newTestKbaseUserFederation(t, filePath) - users, err := kbaseFed.readUserTable() - assert.Nil(err, "Error reading good_user_table_%d.csv", i) - assert.Equal(len(users), len(goodUserMap[i]), "Incorrect number of users read from good_user_table_%d.csv", i) - for orcid, username := range goodUserMap[i] { - readUsername, found := users[orcid] - assert.True(found, "ORCID %s not found in users from good_user_table_%d.csv", orcid, i) - assert.Equal(username, readUsername, "Incorrect username for ORCID %s in good_user_table_%d.csv", orcid, i) - } - } - for i := range badUserTables { - filePath := filepath.Join(testDataDir, fmt.Sprintf("bad_user_table_%d.csv", i)) - kbaseFed := newTestKbaseUserFederation(t, filePath) - users, err := kbaseFed.readUserTable() - assert.NotNil(err, "No error reading bad_user_table_%d.csv", i) - assert.Nil(users, "Users read from bad_user_table_%d.csv", i) - } - for i := range badCSVFormat { - filePath := filepath.Join(testDataDir, fmt.Sprintf("bad_csv_format_%d.csv", i)) - kbaseFed := newTestKbaseUserFederation(t, filePath) - users, err := kbaseFed.readUserTable() - assert.NotNil(err, "No error reading bad_csv_format_%d.csv", i) - assert.Nil(users, "Users read from bad_csv_format_%d.csv", i) - } - kbaseFed := newTestKbaseUserFederation(t, "non_existent_file.csv") - users, err := kbaseFed.readUserTable() - assert.NotNil(err, "No error reading non_existent_file.csv") - assert.Nil(users, "Users read from non_existent_file.csv") -} - func TestIsUsername(t *testing.T) { assert := assert.New(t) diff --git a/databases/kbase_lakehouse/database.go b/databases/kbase_lakehouse/database.go new file mode 100644 index 00000000..bc284a0a --- /dev/null +++ b/databases/kbase_lakehouse/database.go @@ -0,0 +1,140 @@ +// Copyright (c) 2023 The KBase Project and its Contributors +// Copyright (c) 2023 Cohere Consulting, LLC +// +// Permission is hereby granted, free of charge, to any person obtaining a copy of +// this software and associated documentation files (the "Software"), to deal in +// the Software without restriction, including without limitation the rights to +// use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies +// of the Software, and to permit persons to whom the Software is furnished to do +// so, subject to the following conditions: +// +// The above copyright notice and this permission notice shall be included in all +// copies or substantial portions of the Software. +// +// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +// SOFTWARE. + +package kbase_lakehouse + +import ( + "fmt" + "net/http" + + "github.com/google/uuid" + "github.com/mitchellh/mapstructure" + + "github.com/kbase/dts/databases" + "github.com/kbase/dts/databases/kbase" // for user federation + "github.com/kbase/dts/endpoints" +) + +// file database appropriate for handling KBase searches and transfers +// (implements the databases.Database interface) +type Database struct { + // HTTP client that caches queries + Client http.Client + // Name of Globus/S3 lakehouse endpoint + EndpointName string + // KBase user federation mechanism (reused from legacy KBase) + kbaseFed kbase.KBaseUserFederation +} + +type Config struct { + Endpoint string `yaml:"endpoint"` + kbase.KBaseUserFederationConfig `yaml:",inline" mapstructure:",squash"` +} + +func NewDatabase(conf Config) (databases.Database, error) { + // make sure the endpoint is valid + if !endpoints.EndpointExists(conf.Endpoint) { + return nil, fmt.Errorf("invalid endpoint '%s' in kbase database configuration", conf.Endpoint) + } + db := Database{ + EndpointName: conf.Endpoint, + } + + // FIXME: we reuse legacy KBase's user federation spreadsheet to map ORCIDs to + // FIXME: Lakehouse users. This should be replaced when practical. + var err error + db.kbaseFed, err = kbase.NewKBaseUserFederation(conf.KBaseUserFederationConfig) + if err != nil { + return nil, err + } + err = db.kbaseFed.Start() + if err != nil { + return nil, err + } + return &db, nil +} + +func DatabaseConstructor(conf map[string]any) func() (databases.Database, error) { + return func() (databases.Database, error) { + var kbaseConf Config + if err := mapstructure.Decode(conf, &kbaseConf); err != nil { + return nil, err + } + return NewDatabase(kbaseConf) + } +} + +func (db *Database) SpecificSearchParameters() map[string]any { + return nil +} + +func (db *Database) Search(orcid string, params databases.SearchParameters) (databases.SearchResults, error) { + err := fmt.Errorf("Search not implemented for kbase_lakehouse database") + return databases.SearchResults{}, err +} + +func (db *Database) Descriptors(orcid string, fileIds []string) ([]map[string]any, error) { + err := fmt.Errorf("Descriptors not implemented for kbase_lakehouse database") + return nil, err +} + +func (db *Database) EndpointNames() []string { + return []string{db.EndpointName} +} + +func (db *Database) StageFiles(orcid string, fileIds []string) (uuid.UUID, error) { + err := fmt.Errorf("StageFiles not implemented for kbase_lakehouse database") + return uuid.UUID{}, err +} + +func (db *Database) StagingStatus(id uuid.UUID) (databases.StagingStatus, error) { + err := fmt.Errorf("StagingStatus not implemented for kbase_lakehouse database") + return databases.StagingStatusUnknown, err +} + +func (db *Database) Finalize(orcid string, id uuid.UUID) error { + return nil +} + +func (db *Database) LocalUser(orcid string) (string, error) { + return db.kbaseFed.UsernameForOrcid(orcid) +} + +// NOTE: This method is KBase-specific and not part of the Database interface. +// NOTE: It's here to allow us to hand a KBase user's Globus ID over to the Globus S3 Connector. +func (db *Database) GlobusId(orcid string) (uuid.UUID, error) { + return db.kbaseFed.GlobusIdForOrcid(orcid) +} + +func (db Database) Save() (databases.DatabaseSaveState, error) { + // so far, this database has no internal state + return databases.DatabaseSaveState{ + Name: "kbase_lakehouse", + }, nil +} + +func (db *Database) Load(state databases.DatabaseSaveState) error { + return nil // no internal state +} + +func (db *Database) FinalizeDatabase() error { + return db.kbaseFed.Stop() +} diff --git a/databases/nmdc/database_test.go b/databases/nmdc/database_test.go index 58c7c135..1ebb3ace 100644 --- a/databases/nmdc/database_test.go +++ b/databases/nmdc/database_test.go @@ -49,13 +49,13 @@ endpoints: name: NMDC (NERSC) id: ${DTS_GLOBUS_TEST_ENDPOINT} provider: globus - root: / + base_path: / credential: globus globus-nmdc-emsl: name: NMDC Bulk Data Cache id: ${DTS_GLOBUS_TEST_ENDPOINT} provider: globus - root: / + base_path: / credential: globus globus-jdp: name: Globus NERSC DTN @@ -826,6 +826,8 @@ func TestDescriptors(t *testing.T) { } } +/* FIXME: All this metadata in this test has recently vanished, so it seems like + * FIXME: we'll have to keep chasing records. func TestCreditMetadataForStudy(t *testing.T) { assert := assert.New(t) db := Database{ @@ -925,6 +927,7 @@ func TestCreditMetadataForStudy(t *testing.T) { assert.Equal("United States Department of Energy", credit.Funding[0].Funder.OrganizationName, "Credit metadata first funding source name is incorrect") } +*/ func TestPageNumberAndSize(t *testing.T) { assert := assert.New(t) diff --git a/deployment/dts.yaml b/deployment/dts.yaml index 8a1976bf..0da30e84 100644 --- a/deployment/dts.yaml +++ b/deployment/dts.yaml @@ -54,32 +54,37 @@ endpoints: globus-local: name: DTS Local Endpoint id: ${LOCAL_ENDPOINT_ID} - provider: globus + provider: local credential: globus globus-jdp: name: DTS JGI Share id: ${JDP_ENDPOINT_ID} provider: globus credential: globus - root: /dm_archive + data_path: dm_archive globus-kbase: name: KBase Bulk Share id: ${KBASE_ENDPOINT_ID} provider: globus credential: globus - root: /jeff_cohere + base_path: ${KBASE_ENDPOINT_BASEPATH} + data_path: jeff_cohere + globus-kbase-lakehouse: + name: KBase Data Lakehouse Development Environment + id: ${KBASE_LAKEHOUSE_ENDPOINT_ID} + provider: globus + credential: globus + base_path: ${KBASE_LAKEHOUSE_ENDPOINT_BASEPATH} globus-nmdc-nersc: name: NMDC (NERSC) id: ${NMDC_NERSC_ENDPOINT_ID} provider: globus credential: globus - root: / globus-nmdc-emsl: name: NMDC Bulk Data Cache id: ${NMDC_EMSL_ENDPOINT_ID} provider: globus credential: globus - root: / s3-nasa-power: name: NASA POWER (S3) bucket: nasa-power @@ -97,6 +102,11 @@ databases: # databases between which files can be transferred organization: KBase endpoint: globus-kbase data_directory: /data + kbase_lakehouse: + name: KBase Lakehouse + organization: KBase + endpoint: globus-kbase-lakehouse + data_directory: /data nmdc: name: National Microbiome Data Collaborative organization: LBNL, PNNL, ORNL diff --git a/docs/admin/config.md b/docs/admin/config.md index 33a1c8ea..a0588b34 100644 --- a/docs/admin/config.md +++ b/docs/admin/config.md @@ -74,7 +74,7 @@ section are: development work. The default value is `false`. * `double_check_staging`: an optional parameter that, if set to `true`, performs additional checks for staged files. This parameter can be useful for figuring - out the appropriate `root` for an endpoint. + out the appropriate `base_path` for an endpoint. ## `endpoints` @@ -125,9 +125,11 @@ The fields that define the behavior of each endpoint are: a client * `client_secret`: a string containing a secret corresponding to the ID provided by the `client_id` parameter -* `root`: this optional parameter specifies the root directory used by DTS to - refer to files on the underlying filesystem of the endpoint. If left blank, - the root directory is set to `/`. +* `base_path`: this optional parameter specifies the root directory used by DTS to + refer to the path on the underlying filesystem at which the endpoint sits. If left blank, + `base_path` is set to `/`. +* `data_path`: this optional parameter specifies a path on the endpoint (relative to `base_path`) + at which files of interest to a database sit. ## `databases` diff --git a/docs/admin/deployment.md b/docs/admin/deployment.md index e1e627e4..31c4ce08 100644 --- a/docs/admin/deployment.md +++ b/docs/admin/deployment.md @@ -116,7 +116,7 @@ following files: * `dts.gob` - a file containing information about pending and recently finished file transfers, along with any related database-specific state information * `kbase_user_orcids.csv` - a comma-separated variable file associating ORCID - identifiers with KBase users. This file is a temporary mechanism that allows - the DTS to obtain the username of a KBase user given their ORCID. It is - re-read at the top of the hour, making it easy to replace without restarting - a deployment. + identifiers with KBase users, and optionally with Globus IDs (which are UUIDs). + This file is a temporary mechanism that allows the DTS to obtain the username + and Globus ID for a KBase user given their ORCID. It is re-read at the top of + the hour, making it easy to replace without restarting a deployment. diff --git a/dtstest/dtstest.go b/dtstest/dtstest.go index 210101b9..07b89f35 100644 --- a/dtstest/dtstest.go +++ b/dtstest/dtstest.go @@ -31,6 +31,7 @@ import ( "github.com/google/uuid" + "github.com/kbase/dts/auth" "github.com/kbase/dts/config" "github.com/kbase/dts/databases" "github.com/kbase/dts/endpoints" @@ -99,14 +100,17 @@ type EndpointOptions struct { // This type implements an Endpoint test fixture type Endpoint struct { + Id_ uuid.UUID // database fixture attached to endpoint Database *Database // endpoint testing options Options EndpointOptions // a table of ongoing "file transfers" Xfers map[uuid.UUID]transferInfo - // root path - RootPath string + Paths struct { + Base string + Data string + } // a set of files on this endpoint that have been staged StagedFiles map[string]bool } @@ -117,14 +121,19 @@ type Endpoint struct { func RegisterEndpoint(endpointName string, options EndpointOptions) error { slog.Debug(fmt.Sprintf("Registering test endpoint %s...", endpointName)) newEndpointFunc := func(conf map[string]any) (endpoints.Endpoint, error) { - root, ok := config.Endpoints[endpointName]["root"].(string) + basePath, ok := config.Endpoints[endpointName]["base_path"].(string) if !ok { - root = "/" + basePath = "/" } + dataPath, _ := config.Endpoints[endpointName]["data_path"].(string) return &Endpoint{ - Options: options, - Xfers: make(map[uuid.UUID]transferInfo), - RootPath: root, + Id_: uuid.New(), + Options: options, + Xfers: make(map[uuid.UUID]transferInfo), + Paths: struct{ Base, Data string }{ + Base: basePath, + Data: dataPath, + }, StagedFiles: make(map[string]bool), }, nil } @@ -135,12 +144,28 @@ func RegisterEndpoint(endpointName string, options EndpointOptions) error { return endpoints.RegisterEndpointProvider(provider, newEndpointFunc) } +func (ep *Endpoint) Id() uuid.UUID { + return ep.Id_ +} + func (ep *Endpoint) Provider() string { return "dtstest" } -func (ep *Endpoint) Root() string { - return ep.RootPath +func (ep *Endpoint) BasePath() string { + return ep.Paths.Base +} + +func (ep *Endpoint) DataPath() string { + return ep.Paths.Data +} + +func (ep *Endpoint) ConnectsWith(provіder string) bool { + return provіder == "dtstest" +} + +func (ep *Endpoint) RegisterConnectionCredential(user auth.User, provіder string) error { + return nil } func (ep *Endpoint) FilesStaged(files []map[string]any) (bool, error) { @@ -180,7 +205,7 @@ func (ep *Endpoint) Transfers() ([]uuid.UUID, error) { return xfers, nil } -func (ep *Endpoint) Transfer(dst endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { +func (ep *Endpoint) Transfer(user auth.User, dst endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { xferId := uuid.New() ep.Xfers[xferId] = transferInfo{ Time: time.Now(), @@ -313,17 +338,20 @@ func (db *Database) StageFiles(orcid string, fileIds []string) (uuid.UUID, error func (db *Database) StagingStatus(id uuid.UUID) (databases.StagingStatus, error) { if info, found := db.Staging[id]; found { - endpoint := db.Endpt.(*Endpoint) - if time.Since(info.Time) >= endpoint.Options.StagingDuration { // FIXME: not always so! - // update the staged status on the test endpoint - stagingRequest := db.Staging[id] - for _, fileId := range stagingRequest.FileIds { - endpoint.StagedFiles[fileId] = true + if endpoint, ok := db.Endpt.(*Endpoint); ok { + if time.Since(info.Time) >= endpoint.Options.StagingDuration { // FIXME: not always so! + // update the staged status on the test endpoint + stagingRequest := db.Staging[id] + for _, fileId := range stagingRequest.FileIds { + endpoint.StagedFiles[fileId] = true + } + return databases.StagingStatusSucceeded, nil } - + return databases.StagingStatusActive, nil + } else { + // assume staging succeeded return databases.StagingStatusSucceeded, nil } - return databases.StagingStatusActive, nil } return databases.StagingStatusUnknown, nil } diff --git a/endpoints/endpoints.go b/endpoints/endpoints.go index 4d2b42d5..87181681 100644 --- a/endpoints/endpoints.go +++ b/endpoints/endpoints.go @@ -22,10 +22,13 @@ package endpoints import ( + "fmt" + "log/slog" "sync" "github.com/google/uuid" + "github.com/kbase/dts/auth" "github.com/kbase/dts/config" ) @@ -67,10 +70,19 @@ type TransferStatus struct { // This type represents an endpoint for transferring files. type Endpoint interface { + // Returns the endpoint's unique identifier. + Id() uuid.UUID // Returns a string indicating the service provider for the endpoint. Provider() string - // Returns the path on the file system that serves as the endpoint's root. - Root() string + // Returns the path on the file system that serves as the endpoint's base path, below which + // no files are visible. + BasePath() string + // Returns the path of the file system at which files of interest sit (relative to the base path). + // If blank, BasePath is used to locate files. + DataPath() string + // Returns true if this endpoint can transfer files to an endpoint with the given provider, + // false otherwise. + ConnectsWith(provider string) bool // Returns true if the files associated with the given Frictionless // descriptors are staged at this endpoint AND are valid, false otherwise. FilesStaged(descriptors []map[string]any) (bool, error) @@ -78,8 +90,9 @@ type Endpoint interface { Transfers() ([]uuid.UUID, error) // Begins a transfer task that moves the files identified by the FileTransfer // structs, returning a UUID that can be used to refer to this task. It is assumed that there - // no duplicates in the list of files to be transfered. - Transfer(dst Endpoint, files []FileTransfer) (uuid.UUID, error) + // no duplicates in the list of files to be transfered. If authorization is not required for the + // transfer (e.g. DTS performs the transfer on a user's behalf), `user` can be zero-initialized. + Transfer(user auth.User, dst Endpoint, files []FileTransfer) (uuid.UUID, error) // Retrieves the status for a transfer task identified by its UUID. Status(id uuid.UUID) (TransferStatus, error) // Cancels the transfer task with the given UUID (must return immediately, @@ -135,6 +148,11 @@ func NewEndpoint(endpointName string) (Endpoint, error) { } if createEp, valid := createEndpointFuncs_[provider]; valid { endpoint, err = createEp(epConfig) + if err != nil { + return endpoint, err + } + slog.Debug(fmt.Sprintf("Endpoint %s: base path is %s", endpointName, endpoint.BasePath())) + slog.Debug(fmt.Sprintf("Endpoint %s: relative data path is %s", endpointName, endpoint.DataPath())) } else { // invalid provider! err = InvalidProviderError{ Name: endpointName, diff --git a/endpoints/globus/endpoint.go b/endpoints/globus/endpoint.go index f47cdfab..57de546c 100644 --- a/endpoints/globus/endpoint.go +++ b/endpoints/globus/endpoint.go @@ -22,17 +22,12 @@ package globus import ( - "bytes" "encoding/json" - "errors" "fmt" "io" "log/slog" - "net/http" - "net/url" "path/filepath" "strings" - "time" "github.com/google/uuid" "github.com/mitchellh/mapstructure" @@ -44,22 +39,13 @@ import ( // This file implements a Globus endpoint. It uses the Globus Transfer API // described at https://docs.globus.org/api/transfer/. -const ( - globusTransferBaseURL = "https://transfer.api.globusonline.org" - globusTransferApiVersion = "v0.10" -) - -// this error type is returned when a Globus operation fails for any reason -type GlobusError struct { - Code string `json:"code"` - Message string `json:"message"` - - // ConsentRequired error field - RequiredScopes []string `json:"required_scopes"` -} - -func (e GlobusError) Error() string { - return fmt.Sprintf("%s (%s)", e.Message, e.Code) +type GlobusUserCredential struct { + // Authenticated DTS user for whom Globus credential is (temporarily) registered + User auth.User + // (S3) Bucket associated with user transfer + Bucket string + // Globus unique credential identifier + Id uuid.UUID } // this type satisfies the endpoints.Endpoint interface for Globus endpoints @@ -67,18 +53,17 @@ type Endpoint struct { // descriptive endpoint name (obtained from config) Name string // endpoint UUID (obtained from config) - Id uuid.UUID - // root directory for endpoint - RootDir string - // OAuth2 access token - AccessToken string + Id_ uuid.UUID + // Globus clients + Globus GlobusTransferClient + GCSM *GlobusConnectServerManagerClient - // authentication stuff - ClientId uuid.UUID - ClientSecret string + Paths struct { + Base string + Data string + } - // endpoint configuration - Info EndpointInfo + provider string } // configuration struct for Globus endpoints @@ -86,50 +71,45 @@ type Config struct { Name string `yaml:"name"` Id string `yaml:"id"` Credential auth.Credential `yaml:"credential"` - Root string `yaml:"root,omitempty"` + BasePath string `yaml:"base_path,omitempty" mapstructure:"base_path,omitempty"` + DataPath string `yaml:"data_path,omitempty" mapstructure:"data_path,omitempty"` } // creates a new Globus endpoint using the given information func NewEndpoint(config Config) (endpoints.Endpoint, error) { - clientId, err := uuid.Parse(config.Credential.Id) - if err != nil { - return nil, fmt.Errorf("invalid Globus client ID for credential '%s': %s (must be UUID)", - config.Name, config.Credential.Id) - } id, err := uuid.Parse(config.Id) if err != nil { return nil, fmt.Errorf("invalid UUID specified for Globus endpoint: %s", config.Id) } - ep := &Endpoint{ - Name: config.Name, - Id: id, - ClientId: clientId, - ClientSecret: config.Credential.Secret, + globus, err := NewGlobusTransferClient(config.Credential, id) + if err != nil { + return nil, err } - - // if needed, authenticate to obtain a Globus Transfer API access token - var zeroId uuid.UUID - if ep.ClientId != zeroId { - err := ep.authenticate(defaultScopes_) - if err != nil { - return ep, err - } + ep := &Endpoint{ + Name: config.Name, + Id_: id, + Globus: globus, } - // if present, the root entry overrides the endpoint's root, and is expressed - // as a path relative to it - if config.Root != "" { - ep.RootDir = config.Root + if config.BasePath != "" { + ep.Paths.Base = config.BasePath } else { - ep.RootDir = "/" + ep.Paths.Base = "/" } - slog.Debug(fmt.Sprintf("Endpoint %s: root directory is %s", - ep.Name, ep.RootDir)) + ep.Paths.Data = config.DataPath - // query the endpoint for its capabilities - ep.Info, err = ep.getEndpointInfo(ep.Id) + // try accessing the Connect Server Manager API + if ep.Globus.Info.EntityType == "GCSv5_mapped_collection" && ep.Globus.Info.GCSManagerUrl != "" { + ep.GCSM, _ = ep.Globus.ConnectServerManagerClient() + if ep.GCSM != nil { + slog.Debug("Connected to Globus Connect Server Manager.") + } + } - return ep, err + if ep.provider, err = ep.determineProvider(); err != nil { + return nil, err + } + return ep, nil } // constructs a Globus endpoint from a configuration map @@ -142,12 +122,30 @@ func EndpointConstructor(conf map[string]any) (endpoints.Endpoint, error) { return NewEndpoint(globusConfig) } -func (ep *Endpoint) Provider() string { - return "globus" +func (ep Endpoint) Id() uuid.UUID { + return ep.Id_ } -func (ep *Endpoint) Root() string { - return ep.RootDir +func (ep Endpoint) Provider() string { + // A Globus endpoint can have a different provider via Globus Premium Connectors. + return ep.provider +} + +func (ep Endpoint) BasePath() string { + return ep.Paths.Base +} + +func (ep Endpoint) DataPath() string { + return ep.Paths.Data +} + +func (ep Endpoint) ConnectsWith(provider string) bool { + switch provider { + case "globus", "s3": + return true + default: + return false + } } func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { @@ -155,7 +153,7 @@ func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { filesInDir := make(map[string][]string) for _, descriptor := range descriptors { dir, file := filepath.Split(descriptor["path"].(string)) - dir = filepath.Join(ep.RootDir, dir) + dir = filepath.Join(ep.DataPath(), dir) if _, found := filesInDir[dir]; !found { filesInDir[dir] = make([]string, 0) } @@ -163,16 +161,11 @@ func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { } // for each directory, check for its existence and that its files are present - // (https://docs.globus.org/api/transfer/file_operations/#list_directory_contents) for dir, files := range filesInDir { - values := url.Values{} - values.Add("path", dir) - values.Add("orderby", "name ASC") - resource := fmt.Sprintf("operation/endpoint/%s/ls", ep.Id.String()) - body, err := ep.get(resource, values) + globusFiles, err := ep.Globus.FilesInDirectory(dir) if err != nil { switch lsErr := err.(type) { - case *GlobusError: + case *GlobusTransferError: switch lsErr.Code { case "ClientError.NotFound": // it's okay if the directory doesn't exist -- it might need to be staged @@ -186,21 +179,9 @@ func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { return false, err } } - - // https://docs.globus.org/api/transfer/file_operations/#dir_listing_response - type DirListingResponse struct { - Data []struct { - Name string `json:"name"` - } `json:"DATA"` - } - var response DirListingResponse - err = json.Unmarshal(body, &response) - if err != nil { - return false, err - } filesPresent := make(map[string]bool) - for _, data := range response.Data { - filesPresent[data.Name] = true + for _, file := range globusFiles { + filesPresent[file] = true } for _, file := range files { if _, present := filesPresent[file]; !present { @@ -212,59 +193,51 @@ func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { } func (ep *Endpoint) Transfers() ([]uuid.UUID, error) { - // https://docs.globus.org/api/transfer/task/#get_task_list - values := url.Values{} - values.Add("fields", "task_id") - values.Add("filter", "status:ACTIVE,INACTIVE/label:DTS") - values.Add("limit", "1000") - values.Add("orderby", "name ASC") - - body, err := ep.get("task_list", url.Values{}) - if err != nil { - return nil, err - } - type TaskListResponse struct { - Length int `json:"length"` - Limit int `json:"limіt"` - Data []struct { - TaskId uuid.UUID `json:"task_id"` - } `json:"DATA"` - } - var response TaskListResponse - err = json.Unmarshal(body, &response) - if err != nil { - return nil, err - } - taskIds := make([]uuid.UUID, len(response.Data)) - for i, data := range response.Data { - taskIds[i] = data.TaskId - } - return taskIds, nil + return ep.Globus.TransferTasks() } -func (ep *Endpoint) Transfer(destination endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { +func (ep Endpoint) Transfer(user auth.User, destination endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { + if _, isGlobus := destination.(*Endpoint); !isGlobus { + return uuid.UUID{}, &endpoints.IncompatibleDestinationError{ + Source: ep.Id().String(), + SourceProvider: ep.Provider(), + Destination: destination.Id().String(), + DestinationProvider: destination.Provider(), + Message: "a premium Globus connector may be required", + } + } + // NOTE: We don't check whether files are staged here, because the endpoint itself doesn't always // have a reliable staging check (e.g. JDP's private data is invisible to Globus directory // listings). Consequently, we assume that files are staged by the time this function is called. - // obtain a submission ID - submissionId, err := ep.getSubmissionId() - if err != nil { - return uuid.UUID{}, err + filesWithFullPath := make([]endpoints.FileTransfer, len(files)) + for i, file := range files { + filesWithFullPath[i] = endpoints.FileTransfer{ + SourcePath: filepath.Join(ep.DataPath(), file.SourcePath), + DestinationPath: file.DestinationPath, + Hash: file.Hash, + HashAlgorithm: file.HashAlgorithm, + } } - // Occasionally, Globus returns a zero-valued UUID (uuid.Nil) and no error (network burp?). - // So we pause and resubmit in this case - for submissionId == uuid.Nil { - time.Sleep(time.Second) - submissionId, err = ep.getSubmissionId() - if err != nil { - return uuid.UUID{}, err + // If this is a transfer between endpoints with different providers, register or fetch the + // credential that allows them to connect. + var credential auth.Credential + if ep.Provider() != destination.Provider() { + slog.Debug("Source and destination providers differ, registering credentials...") + if destEp, isGlobus := destination.(*Endpoint); isGlobus { + if destEp.GCSM == nil { + return uuid.UUID{}, fmt.Errorf("the Globus Connect Server Manager API is not available; cannot register credentials") + } + var err error + if credential, err = destEp.GCSM.AddOrUpdateUserCredential(user, destination.Provider()); err != nil { + return uuid.UUID{}, err + } } } - // now, submit the transfer task itself - return ep.submitTransfer(destination, submissionId, files) + return ep.Globus.Transfer(credential, ep.Id(), destination.Id(), filesWithFullPath) } // mapping of Globus status code strings to DTS status codes @@ -276,49 +249,23 @@ var statusCodesForStrings = map[string]endpoints.TransferStatusCode{ } func (ep *Endpoint) Status(id uuid.UUID) (endpoints.TransferStatus, error) { - resource := fmt.Sprintf("task/%s", id.String()) - body, err := ep.get(resource, url.Values{}) - if err != nil { - return endpoints.TransferStatus{}, err - } - if responseIsError(body) { - var globusErr GlobusError - err := json.Unmarshal(body, &globusErr) - if err == nil { - err = &globusErr - } - return endpoints.TransferStatus{}, err - } - type TaskResponse struct { - Files int `json:"files"` - FilesSkipped int `json:"files_skipped"` - FilesTransferred int `json:"files_transferred"` - IsPaused bool `json:"is_paused"` - NiceStatus string `json:"nice_status"` - NiceStatusShortDescription string `json:"nice_status_short_description"` - Status string `json:"status"` - } - var response TaskResponse - err = json.Unmarshal(body, &response) + taskStatus, err := ep.Globus.TaskStatus(id) if err != nil { return endpoints.TransferStatus{}, err } + // check for an error condition in NiceStatus - if response.NiceStatus != "" && response.NiceStatus != "OK" && response.NiceStatus != "Queued" { + if taskStatus.NiceStatus != "" && taskStatus.NiceStatus != "OK" && taskStatus.NiceStatus != "Queued" { // get the event list for this task - resource := fmt.Sprintf("task/%s/event_list", id.String()) - body, err := ep.get(resource, url.Values{}) + events, err := ep.Globus.TaskEvents(id) if err != nil { - // fine, we'll just use the "nice status" - return endpoints.TransferStatus{}, errors.New(response.NiceStatusShortDescription) + return endpoints.TransferStatus{}, err } - var eventList EventList - json.Unmarshal(body, &eventList) - if response.NiceStatus == "AUTH" { + if taskStatus.NiceStatus == "AUTH" { // sometimes Globus throws an AUTH error here during a network burp, so we // ignore it and report a failed status check (after all, we can't get here // without AUTHing successfully!) - for _, event := range eventList.Data { + for _, event := range events { if event.IsError { slog.Debug(fmt.Sprintf("Globus task %s: status check failed with AUTH error below (probably bogus, ignoring): ", id.String())) slog.Debug(fmt.Sprintf("Globus task %s: %s (%s):\n%s", id.String(), event.Description, event.Code, event.Details)) @@ -328,389 +275,53 @@ func (ep *Endpoint) Status(id uuid.UUID) (endpoints.TransferStatus, error) { // it's probably real, so traverse the event list return endpoints.TransferStatus{ Code: endpoints.TransferStatusFailed, - Message: descriptionFromEventList(eventList, response.NiceStatusShortDescription), - NumFiles: response.Files, - NumFilesSkipped: response.FilesSkipped, - NumFilesTransferred: response.FilesTransferred, + Message: descriptionFromEventList(events, taskStatus.NiceStatusShortDescription), + NumFiles: taskStatus.Files, + NumFilesSkipped: taskStatus.FilesSkipped, + NumFilesTransferred: taskStatus.FilesTransferred, }, nil } } return endpoints.TransferStatus{ - Code: statusCodesForStrings[response.Status], - NumFiles: response.Files, - NumFilesSkipped: response.FilesSkipped, - NumFilesTransferred: response.FilesTransferred, + Code: statusCodesForStrings[taskStatus.Status], + NumFiles: taskStatus.Files, + NumFilesSkipped: taskStatus.FilesSkipped, + NumFilesTransferred: taskStatus.FilesTransferred, }, nil } func (ep *Endpoint) Cancel(id uuid.UUID) error { - // Because cancellation requests can't be honored under all circumstances, - // this Globus call is asynchronous. Nevertheless, the Globus documentation - // (https://docs.globus.org/api/transfer/task/#cancel_task_by_id) claims the - // call can take up to 10 seconds before returning, which doesn't meet the - // needs of the DTS. The possible outcomes of the call are identified with - // these response codes: - // 1. "Canceled", indicating that the task has been canceled - // 2. "CancelAccepted", indicating that the cancellation request has been - // acknowledged but not yet processed - // 3. "TaskComplete", indicating that the task is complete and not able to - // be canceled. - // - // We live with the 10-second wait for now, since our polling interval is - // large. - resource := fmt.Sprintf("task/%s/cancel", id.String()) - _, err := ep.post(resource, nil) // can take up to 10 ѕeconds! - // FIXME: if this ^^^ becomes an issue, we can dispatch the POST to a - // FIXME: persistent goroutine to handle the cancellation - if err != nil { - if globusError, ok := err.(*GlobusError); ok { - switch globusError.Code { - case "Canceled", "CancelAccepted", "TaskComplete": // it worked! - err = nil - } - } - } - return err + return ep.Globus.Cancel(id) } -//----------- -// Internals -//----------- - -// default client credentials grant scopes -var defaultScopes_ = []string{"urn:globus:auth:scope:transfer.api.globus.org:all"} - -// returns true if a Globus response body matches an error -func responseIsError(body []byte) bool { - bodyStr := string(body) - return strings.Contains(bodyStr, "\"code\"") && - !strings.Contains(bodyStr, "\"code\": \"Accepted\"") && - strings.Contains(string(body), "\"message\"") -} - -// (re)authenticates with Globus using its client ID and secret to obtain an -// access token with consents for its relevant list of scopes -// (https://docs.globus.org/api/auth/reference/#client_credentials_grant) -func (ep *Endpoint) authenticate(scopes []string) error { - authUrl := "https://auth.globus.org/v2/oauth2/token" - data := url.Values{} - data.Set("scope", strings.Join(scopes, " ")) - data.Set("grant_type", "client_credentials") - req, err := http.NewRequest(http.MethodPost, authUrl, strings.NewReader(data.Encode())) - if err != nil { - return err - } - req.SetBasicAuth(ep.ClientId.String(), ep.ClientSecret) - req.Header.Add("Content-Type", "application-x-www-form-urlencoded") - - // send the request using a fresh HTTP client - var client http.Client - resp, err := client.Do(req) - if err != nil { - return err - } - if resp.StatusCode != 200 { - // fish specifics out of the response - type AuthError struct { - Error string `json:"error"` - Description string `json:"error_description"` - URI string `json:"error_uri"` - } - body, err := io.ReadAll(resp.Body) - if err != nil { - return err - } - var authError AuthError - err = json.Unmarshal(body, &authError) - if err != nil { - // report the authentication error without details - return fmt.Errorf("couldn't authenticate via Globus Auth API (%d)", resp.StatusCode) - } - if len(authError.Description) > 0 { - return fmt.Errorf("couldn't authenticate via Globus Auth API: %s; %s (%d)", - authError.Error, authError.Description, resp.StatusCode) - } - return fmt.Errorf("couldn't authenticate via Globus Auth API: %s (%d)", - authError.Error, resp.StatusCode) - } - - // read and unmarshal the response - body, err := io.ReadAll(resp.Body) - if err != nil { - return err - } - type AuthResponse struct { - AccessToken string `json:"access_token"` - Scope string `json:"scope"` - ResourceServer string `json:"resource_server"` - ExpiresIn int `json:"expires_in"` - TokenType string `json:"token_type"` - } - var authResponse AuthResponse - err = json.Unmarshal(body, &authResponse) +// Performs an HTTPS PUT request on the endpoint, uploading the content of the given reader as +// the request body. Only supported if the Globus endpoint has an associated HTTPS server. +func (ep *Endpoint) PutFromReader(resource string, body io.Reader) error { + httpsClient, err := ep.Globus.HttpsClient(ep.Id()) if err != nil { return err } - - // FIXME: check the scopes to see if they match our requested ones? - - // stash the access token - ep.AccessToken = authResponse.AccessToken - - return nil + absPath := filepath.Join(ep.Paths.Base, resource) + return httpsClient.PutFile(absPath, body) } -// This helper sends the given HTTP request, parsing the response for -// Globus-style error codes/messages and handling the ones that can be -// handled automatically (e.g. consent/scope related errors). In any case, -// it returns a byte slice containing the body of the response or an -// error indicating failure. -func (ep *Endpoint) sendRequest(request *http.Request) ([]byte, error) { - // send the initial request with a fresh HTTP client - var client http.Client - resp, err := client.Do(request) - if err != nil { - return nil, err - } - body, err := io.ReadAll(resp.Body) - if err != nil { - return nil, err - } - resp.Body.Close() +//----------- +// Internals +//----------- - // check the response for a Globus-style error code / message - if responseIsError(body) { - var errResp GlobusError - err = json.Unmarshal(body, &errResp) - if err != nil { - return nil, err - } - if errResp.Code == "ConsentRequired" || errResp.Code == "AuthenticationFailed" { - // our token has expired or we're missing a required scope, - // so reauthenticate - if len(errResp.RequiredScopes) > 0 { - err = ep.authenticate(errResp.RequiredScopes) - } else { - err = ep.authenticate(defaultScopes_) - } - if err != nil { - return nil, err +func (ep Endpoint) determineProvider() (string, error) { + if ep.GCSM != nil { + // sift through the storage providers in the gateways + // NOTE: we match the first policy we find + for _, gateway := range ep.GCSM.StorageGateways { + slog.Debug(fmt.Sprintf("Storage gateway provider: %s", gateway.Provider)) + if gateway.Provider == "s3" { + return "s3", nil } - // try the request again - resp, err = client.Do(request) - if err != nil { - return nil, err - } - body, err = io.ReadAll(resp.Body) - resp.Body.Close() - } else { - // other errors are propagated - err = &errResp - } - } - return body, err -} - -// Performs a GET request on the given Globus resource, handling any obvious -// errors and returning a byte slice containing the body of the response, -// and/or any unhandled error. -// This method handles scope-related errors by reauthenticating as needed and -// retrying the operation. See https://docs.globus.org/api/flows/working-with-consents/ -// for details on Globus scopes and consents. -func (ep *Endpoint) get(resource string, values url.Values) ([]byte, error) { - u, err := url.ParseRequestURI(globusTransferBaseURL) - if err != nil { - return nil, err - } - u.Path = fmt.Sprintf("%s/%s", globusTransferApiVersion, resource) - u.RawQuery = values.Encode() - res := fmt.Sprintf("%v", u) - slog.Debug(fmt.Sprintf("GET: %s", res)) - req, err := http.NewRequest(http.MethodGet, res, http.NoBody) - if err != nil { - return nil, err - } - req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", ep.AccessToken)) - - return ep.sendRequest(req) -} - -// Performs a POST request on the given Globus resource, handling any obvious -// errors and returning a byte slice containing the body of the response, -// and/or any unhandled error. -// This method handles scope-related errors by reauthenticating as needed and -// retrying the operation. See https://docs.globus.org/api/flows/working-with-consents/ -// for details on Globus scopes and consents. -func (ep *Endpoint) post(resource string, body io.Reader) ([]byte, error) { - u, err := url.ParseRequestURI(globusTransferBaseURL) - if err != nil { - return nil, err - } - u.Path = fmt.Sprintf("%s/%s", globusTransferApiVersion, resource) - res := fmt.Sprintf("%v", u) - slog.Debug(fmt.Sprintf("POST: %s", res)) - req, err := http.NewRequest(http.MethodPost, res, body) - if err != nil { - return nil, err - } - req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", ep.AccessToken)) - req.Header.Set("Content-Type", "application/json") - - return ep.sendRequest(req) -} - -// https://docs.globus.org/api/transfer/task_submit/#get_submission_id -func (ep *Endpoint) getSubmissionId() (uuid.UUID, error) { - var id uuid.UUID - body, err := ep.get("submission_id", url.Values{}) - if err != nil { - return id, err - } - type SubmissionIdResponse struct { - Value uuid.UUID `json:"value"` - } - var response SubmissionIdResponse - err = json.Unmarshal(body, &response) - return response.Value, err -} - -// https://docs.globus.org/api/transfer/endpoints_and_collections/#get_endpoint_or_collection_by_id -// https://docs.globus.org/api/transfer/task_submit/#submit_transfer_task -// https://docs.globus.org/api/transfer/task_submit/#transfer_item_fields -func (ep *Endpoint) submitTransfer(destination endpoints.Endpoint, - submissionId uuid.UUID, files []endpoints.FileTransfer) (uuid.UUID, error) { - var xferId uuid.UUID - - // are the source and destination endpoints configured in a conflicting way? - globusDestination := destination.(*Endpoint) - if ep.Info.ForceVerify && globusDestination.Info.DisableVerify { // not allowed! - return xferId, &endpoints.IncompatibleDestinationError{ - Source: ep.Name, - SourceProvider: "globus", - Destination: globusDestination.Name, - DestinationProvider: "globus", - Message: "Source endpoint forces checksum verification, but destination disables it.", - } - } - - // configure checksum settings based on destination endpoint info - var verifyChecksum bool = true - var syncLevel int = 3 // transfer only if checksums don't match - if globusDestination.Info.DisableVerify { // checksum verification disabled on endpoint - verifyChecksum = false - syncLevel = 2 // transfer if source file is newer than destination file - } - - type TransferItem struct { - DataType string `json:"DATA_TYPE"` // "transfer_item" - SourcePath string `json:"source_path"` - DestinationPath string `json:"destination_path"` - ExternalChecksum string `json:"external_checksum,omitempty"` - ChecksumAlgorithm string `json:"checksum_algorithm,omitempty"` - } - xferItems := make([]TransferItem, len(files)) - for i, file := range files { - var checksum, checksumAlgorithm string - if verifyChecksum { - checksum = file.Hash - checksumAlgorithm = file.HashAlgorithm - } - xferItems[i] = TransferItem{ - DataType: "transfer_item", - SourcePath: filepath.Join(ep.RootDir, file.SourcePath), - DestinationPath: file.DestinationPath, - ExternalChecksum: checksum, - ChecksumAlgorithm: checksumAlgorithm, - } - } - - // the destination is a Globus endpoint, right? - gDestination, ok := destination.(*Endpoint) - if !ok { - return xferId, &endpoints.IncompatibleDestinationError{ - Source: ep.Name, - SourceProvider: "globus", - Destination: "???", - DestinationProvider: destination.Provider(), - Message: "destination is not a Globus endpoint", - } - } - - // submit the transfer request - type SubmissionRequest struct { - DataType string `json:"DATA_TYPE"` // "transfer" - Id string `json:"submission_id"` - Label string `json:"label"` // "DTS" - Data []TransferItem `json:"DATA"` - DestinationEndpoint string `json:"destination_endpoint"` - SourceEndpoint string `json:"source_endpoint"` - SyncLevel int `json:"sync_level"` - VerifyChecksum bool `json:"verify_checksum"` - FailOnQuotaErrors bool `json:"fail_on_quota_errors"` - } - data, err := json.Marshal(SubmissionRequest{ - DataType: "transfer", - Id: submissionId.String(), - Label: "DTS", - Data: xferItems, - DestinationEndpoint: gDestination.Id.String(), - SourceEndpoint: ep.Id.String(), - SyncLevel: syncLevel, - VerifyChecksum: verifyChecksum, - FailOnQuotaErrors: true, - }) - if err != nil { - return xferId, err - } - body, err := ep.post("transfer", bytes.NewReader(data)) - if err != nil { - return xferId, err - } - if responseIsError(body) { - var globusErr GlobusError - err = json.Unmarshal(body, &globusErr) - if err == nil { - err = &globusErr - } - return xferId, err - } - type SubmissionResponse struct { - TaskId uuid.UUID `json:"task_id"` - } - - var gResp SubmissionResponse - err = json.Unmarshal(body, &gResp) - if err != nil { - return xferId, err - } - xferId = gResp.TaskId - slog.Debug(fmt.Sprintf("Initiated Globus transfer task %s (%d files)", - xferId.String(), len(files))) - return xferId, nil -} - -type EndpointInfo struct { - DisableVerify bool `json:"disable_verify"` // true if checksums are not available - ForceVerify bool `json:"force_verify"` // true if checksums must be available -} - -func (ep *Endpoint) getEndpointInfo(id uuid.UUID) (EndpointInfo, error) { - // query the endpoint for its capabilities - body, err := ep.get(fmt.Sprintf("endpoint/%s", id), url.Values{}) - if err != nil { - return EndpointInfo{}, err - } - if responseIsError(body) { - var globusErr GlobusError - err = json.Unmarshal(body, &globusErr) - if err == nil { - err = &globusErr } - return EndpointInfo{}, err + return "globus", nil } - var endpointInfo EndpointInfo - err = json.Unmarshal(body, &endpointInfo) - return endpointInfo, err + return "globus", nil } type EventList struct { @@ -728,10 +339,10 @@ type Event struct { // traverses a Globus event list, producing an appropriate description of errors encountered, // falling back to the given description if nothing can be gleaned -func descriptionFromEventList(events EventList, fallback string) string { +func descriptionFromEventList(events []GlobusEvent, fallback string) string { missing_files := make(map[string]bool) inaccessible_files := make(map[string]bool) - for _, event := range events.Data { + for _, event := range events { if event.IsError { switch event.Code { case "FILE_NOT_FOUND", "PERMISSION_DENIED": diff --git a/endpoints/globus/endpoint_test.go b/endpoints/globus/endpoint_test.go index c4c5185f..855b0476 100644 --- a/endpoints/globus/endpoint_test.go +++ b/endpoints/globus/endpoint_test.go @@ -34,6 +34,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert/yaml" + "github.com/kbase/dts/auth" "github.com/kbase/dts/endpoints" ) @@ -179,12 +180,13 @@ func TestGlobusConstructor(t *testing.T) { assert.Nil(err) endpoint, err := EndpointConstructor(configMap) - assert.NotNil(endpoint) // if invalid credientials are provided, an error is returned if !checkGlobusEnvVars() { + assert.Nil(endpoint) assert.NotNil(err) return } + assert.NotNil(endpoint) assert.Nil(err) } @@ -309,7 +311,7 @@ func TestGlobusTransfer(t *testing.T) { DestinationPath: path.Join(destDirName(16), path.Base(sourceFilesById[id])), }) } - taskId, err := source.Transfer(destination, fileXfers) + taskId, err := source.Transfer(auth.User{}, destination, fileXfers) assert.Nil(err) // wait for the task to register in the system @@ -388,7 +390,7 @@ func TestGlobusTransferCancellation(t *testing.T) { DestinationPath: path.Join(destDirName(16), path.Base(sourceFilesById[id])), }) } - taskId, err := source.Transfer(destination, fileXfers) + taskId, err := source.Transfer(auth.User{}, destination, fileXfers) assert.Nil(err) // wait for the task to show up diff --git a/endpoints/globus/globus.go b/endpoints/globus/globus.go new file mode 100644 index 00000000..47fae4ac --- /dev/null +++ b/endpoints/globus/globus.go @@ -0,0 +1,1081 @@ +// Copyright (c) 2023 The KBase Project and its Contributors +// Copyright (c) 2023 Cohere Consulting, LLC +// +// Permission is hereby granted, free of charge, to any person obtaining a copy of +// this software and associated documentation files (the "Software"), to deal in +// the Software without restriction, including without limitation the rights to +// use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies +// of the Software, and to permit persons to whom the Software is furnished to do +// so, subject to the following conditions: +// +// The above copyright notice and this permission notice shall be included in all +// copies or substantial portions of the Software. +// +// THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR +// IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, +// FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE +// AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER +// LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, +// OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE +// SOFTWARE. + +package globus + +import ( + "bytes" + "encoding/json" + "errors" + "fmt" + "io" + "log/slog" + "net/http" + "net/url" + "strings" + "time" + + "github.com/google/uuid" + + "github.com/kbase/dts/auth" + "github.com/kbase/dts/endpoints" +) + +// This file implements a Globus endpoint. It uses the Globus Transfer API +// described at https://docs.globus.org/api/transfer/. + +// this error type is returned when a Globus transfer operation fails for any reason +type GlobusTransferError struct { + Code string `json:"code"` + Message string `json:"message"` + + // ConsentRequired error field + RequiredScopes []string `json:"required_scopes"` +} + +func (e GlobusTransferError) Error() string { + return fmt.Sprintf("%s (%s)", e.Message, e.Code) +} + +// this error type is returned when a non-transfer Globus operation fails for any reason +type GlobusGenericError struct { + Message string +} + +func (e GlobusGenericError) Error() string { + return e.Message +} + +// this error type encodes authentication errors and diagnostics +type GlobusAuthRequirementsError struct { + AuthorizationParameters GlobusAuthorizationParameters + Message string + Code string +} + +func (e GlobusAuthRequirementsError) Error() string { + s := fmt.Sprintf("%s (%s)", e.Message, e.Code) + if e.AuthorizationParameters.SessionMessage != "" { + s += ": " + e.AuthorizationParameters.SessionMessage + } + if e.AuthorizationParameters.SessionRequiredIdentities != nil { + s += fmt.Sprintf("; required identities: %v", e.AuthorizationParameters.SessionRequiredIdentities) + } + if e.AuthorizationParameters.SessionRequiredPolicies != nil { + s += fmt.Sprintf("; required policies: %v", e.AuthorizationParameters.SessionRequiredPolicies) + } + if e.AuthorizationParameters.SessionRequiredSingleDomain != nil { + s += fmt.Sprintf("; required identities: %v", e.AuthorizationParameters.SessionRequiredSingleDomain) + } + if e.AuthorizationParameters.SessionRequiredMfa { + s += "; MFA required" + } + if e.AuthorizationParameters.RequiredScopes != nil { + s += fmt.Sprintf("; required scopes: %v", e.AuthorizationParameters.RequiredScopes) + } + if e.AuthorizationParameters.Prompt != "" { + s += "; prompt: " + e.AuthorizationParameters.Prompt + } + return s +} + +// this error indicates that a Globus endpoint has no associated HTTPS server +type GlobusHttpsClientNotAvailableError struct { + Endpoint uuid.UUID +} + +func (e GlobusHttpsClientNotAvailableError) Error() string { + return fmt.Sprintf("no HTTPS Server is not available for endpoint %s", + e.Endpoint.String()) +} + +// this error indicates that the Globus Connect Server Manager client is not available for the +// endpoint in question +type GlobusConnectServerManagerNotAvailableError struct { + Endpoint uuid.UUID +} + +func (e GlobusConnectServerManagerNotAvailableError) Error() string { + return fmt.Sprintf("the Globus Connect Manager Server API is not available for endpoint %s", + e.Endpoint.String()) +} + +// this error contains information about a failed operation with the Globus Connect Server Manager +type GlobusConnectServerManagerError struct { +} + +type GlobusEndpointInfo struct { + DisableVerify bool `json:"disable_verify"` // true if checksums are not available + EntityType string `json:"entity_type"` // indicates type of Globus endpoint server + ForceVerify bool `json:"force_verify"` // true if checksums must be available + GCSManagerUrl string `json:"gcs_manager_url"` // non-blank if GCS Manager operations are supported + HighAssurance bool `json:"high_assurance"` // true if endpoint is a connector + HttpsServer string `json:"https_server"` // non-blank if HTTPS transfers are supported + MappedCollectionId string `json:"mapped_collection_id"` // non-blank if GCS Manager operations are supported + NonfunctionalEndpointId string `json:"non_functional_endpoint_id"` +} + +type GlobusTransferStatus struct { + Files int `json:"files"` + FilesSkipped int `json:"files_skipped"` + FilesTransferred int `json:"files_transferred"` + IsPaused bool `json:"is_paused"` + NiceStatus string `json:"nice_status"` + NiceStatusShortDescription string `json:"nice_status_short_description"` + Status string `json:"status"` +} + +// Globus Transfer API +// https://docs.globus.org/api/transfer/ +type GlobusTransferClient struct { + AccessToken string + Auth *GlobusAuthClient + Scopes []string + EndpointId uuid.UUID + Info GlobusEndpointInfo +} + +// Globus Auth API +// https://docs.globus.org/api/auth/ +type GlobusAuthClient struct { + Credential auth.Credential + Url string +} + +// Globus HTTPS upload client +// https://docs.globus.org/globus-connect-server/v5/https-access-collections +type GlobusHttpsClient struct { + AccessToken string + Scopes []string + Url, DataPath string +} + +type GlobusStorageGateway struct { + ConnectorId uuid.UUID + Id uuid.UUID + Provider string // "s3", etc +} + +// Globus Connect Server Manager API +// https://docs.globus.org/globus-connect-server/v5.4/api/ +type GlobusConnectServerManagerClient struct { + AccessToken string + ClientId string // credential ID that granted access token + EndpointId uuid.UUID + Scopes []string + Url string + StorageGateways []GlobusStorageGateway +} + +// Auth error diagnostics (can be encoded in Globus service responses) +// https://docs.globus.org/guides/overviews/gares/ +type GlobusAuthorizationParameters struct { + SessionMessage string `json:"session_message,omitempty"` + SessionRequiredIdentities []string `json:"session_required_identities,omitempty"` + SessionRequiredPolicies []string `json:"session_required_policies,omitempty"` + SessionRequiredSingleDomain []string `json:"session_required_single_domain,omitempty"` + SessionRequiredMfa bool `json:"session_required_mfa,omitempty"` + RequiredScopes []string `json:"required_scopes,omitempty"` + Prompt string `json:"prompt,omitempty"` +} + +func NewGlobusTransferClient(credential auth.Credential, endpointId uuid.UUID) (GlobusTransferClient, error) { + auth, err := NewGlobusAuthClient(credential) + if err != nil { + return GlobusTransferClient{}, err + } + t := GlobusTransferClient{ + Auth: auth, + EndpointId: endpointId, + Scopes: []string{ + "urn:globus:auth:scope:transfer.api.globus.org:all", + }, + } + if t.AccessToken, err = t.Auth.Authenticate(t.Scopes); err != nil { + return GlobusTransferClient{}, err + } + if t.Info, err = t.getEndpointInfo(endpointId); err != nil { + return GlobusTransferClient{}, err + } + return t, nil +} + +func (t GlobusTransferClient) HttpsClient(endpointId uuid.UUID) (GlobusHttpsClient, error) { + if t.Info.HttpsServer == "" { + return GlobusHttpsClient{}, &GlobusHttpsClientNotAvailableError{t.EndpointId} + } + h := GlobusHttpsClient{ + Scopes: []string{ + fmt.Sprintf("https://auth.globus.org/scopes/%s/https", endpointId.String()), + }, + Url: t.Info.HttpsServer, + } + var err error + if h.AccessToken, err = t.Auth.Authenticate(h.Scopes); err != nil { + return GlobusHttpsClient{}, err + } + return h, nil +} + +func (t GlobusTransferClient) ConnectServerManagerClient() (*GlobusConnectServerManagerClient, error) { + if t.Info.GCSManagerUrl == "" { + return nil, &GlobusConnectServerManagerNotAvailableError{Endpoint: t.EndpointId} + } + scopes := []string{fmt.Sprintf("urn:globus:auth:scope:%s:manage_collections", t.Info.NonfunctionalEndpointId)} + //if !t.Info.HighAssurance { + // scopes[0] += fmt.Sprintf("[*:https://auth.globus.org/scopes/%s/data_access]", t.EndpointId) + //} + m := &GlobusConnectServerManagerClient{ + ClientId: t.Auth.Credential.Id, + EndpointId: t.EndpointId, + Scopes: scopes, + Url: t.Info.GCSManagerUrl, + } + var err error + if m.AccessToken, err = t.Auth.Authenticate(m.Scopes); err != nil { + return nil, err + } + + err = m.getStorageGatewayInfo() + return m, err +} + +// creates a new Globus endpoint using the given information +func NewGlobusAuthClient(credential auth.Credential) (*GlobusAuthClient, error) { + return &GlobusAuthClient{ + Credential: credential, + Url: "https://auth.globus.org/v2/oauth2/token", + }, nil +} + +// (re)authenticates with Globus using its client ID and secret to obtain an +// access token with consents for its relevant list of scopes +// (https://docs.globus.org/api/auth/reference/#client_credentials_grant) +// returns an access token corresponding to the given set of scopes +func (c GlobusAuthClient) Authenticate(scopes []string) (string, error) { + data := url.Values{} + data.Set("scope", strings.Join(scopes, " ")) + data.Set("grant_type", "client_credentials") + req, err := http.NewRequest(http.MethodPost, c.Url, strings.NewReader(data.Encode())) + if err != nil { + return "", err + } + req.SetBasicAuth(c.Credential.Id, c.Credential.Secret) + req.Header.Add("Content-Type", "application-x-www-form-urlencoded") + + // send the request using a fresh HTTP client + var client http.Client + resp, err := client.Do(req) + if err != nil { + return "", err + } + if resp.StatusCode != 200 { + // fish specifics out of the response + type AuthError struct { + Error string `json:"error"` + Description string `json:"error_description"` + URI string `json:"error_uri"` + } + body, err := io.ReadAll(resp.Body) + if err != nil { + return "", err + } + var authError AuthError + err = json.Unmarshal(body, &authError) + if err != nil { + // report the authentication error without details + return "", fmt.Errorf("couldn't authenticate via Globus Auth API (%d)", resp.StatusCode) + } + switch authError.Error { + case "unknown_scope_error": + return "", fmt.Errorf("couldn't authenticate via Globus Auth API: unknown scope(s) requested: %v (%d)", + scopes, resp.StatusCode) + case "invalid_scope_error": + return "", fmt.Errorf("couldn't authenticate via Globus Auth API: invalid scope(s) requested: %v (%d)", + scopes, resp.StatusCode) + default: + if len(authError.Description) > 0 { + return "", fmt.Errorf("couldn't authenticate via Globus Auth API: %s; %s (%d)", + authError.Error, authError.Description, resp.StatusCode) + } + return "", fmt.Errorf("couldn't authenticate via Globus Auth API: %s (%d)", + authError.Error, resp.StatusCode) + } + } + + // read and unmarshal the response + body, err := io.ReadAll(resp.Body) + if err != nil { + return "", err + } + type AuthResponse struct { + AccessToken string `json:"access_token"` + Scope string `json:"scope"` + ResourceServer string `json:"resource_server"` + ExpiresIn int `json:"expires_in"` + TokenType string `json:"token_type"` + } + var authResponse AuthResponse + err = json.Unmarshal(body, &authResponse) + if err != nil { + return "", err + } + + // FIXME: check the scopes to see if they match our requested ones? + + // stash the access token + return authResponse.AccessToken, nil +} + +// https://docs.globus.org/api/transfer/file_operations/#dir_listing_response +// (https://docs.globus.org/api/transfer/file_operations/#list_directory_contents) +func (c *GlobusTransferClient) FilesInDirectory(dir string) ([]string, error) { + values := url.Values{} + values.Add("path", dir) + values.Add("orderby", "name ASC") + body, err := c.get(fmt.Sprintf("operation/endpoint/%s/ls", c.EndpointId), values) + if err != nil { + return nil, err + } + + type DirListingResponse struct { + Data []struct { + Name string `json:"name"` + } `json:"DATA"` + } + var response DirListingResponse + err = json.Unmarshal(body, &response) + if err != nil { + return nil, err + } + files := make([]string, len(response.Data)) + for i, datum := range response.Data { + files[i] = datum.Name + } + return files, nil +} + +func (c *GlobusTransferClient) TransferTasks() ([]uuid.UUID, error) { + // https://docs.globus.org/api/transfer/task/#get_task_list + values := url.Values{} + values.Add("fields", "task_id") + values.Add("filter", "status:ACTIVE,INACTIVE/label:DTS") + values.Add("limit", "1000") + values.Add("orderby", "name ASC") + + body, err := c.get("task_list", url.Values{}) + if err != nil { + return nil, err + } + type TaskListResponse struct { + Length int `json:"length"` + Limit int `json:"limіt"` + Data []struct { + TaskId uuid.UUID `json:"task_id"` + } `json:"DATA"` + } + var response TaskListResponse + err = json.Unmarshal(body, &response) + if err != nil { + return nil, err + } + taskIds := make([]uuid.UUID, len(response.Data)) + for i, data := range response.Data { + taskIds[i] = data.TaskId + } + return taskIds, nil +} + +// Transfers files from the given source endpoint to the given destination endpoint. +// NOTE: file paths are relative to the root of the Globus collection, NOT its +// NOTE: "data directory" +func (c *GlobusTransferClient) Transfer(credential auth.Credential, sourceId, destinationId uuid.UUID, files []endpoints.FileTransfer) (uuid.UUID, error) { + // obtain a submission ID + submissionId, err := c.getSubmissionId() + if err != nil { + return uuid.UUID{}, err + } + + // Occasionally, Globus returns a zero-valued UUID (uuid.Nil) and no error (network burp?). + // So we pause and resubmit in this case + for submissionId == uuid.Nil { + time.Sleep(time.Second) + submissionId, err = c.getSubmissionId() + if err != nil { + return uuid.UUID{}, err + } + } + + // now, submit the transfer task itself + return c.submitTransfer(credential, sourceId, destinationId, submissionId, files) +} + +func (c *GlobusTransferClient) getEndpointInfo(id uuid.UUID) (GlobusEndpointInfo, error) { + // query the endpoint for its capabilities + body, err := c.get(fmt.Sprintf("endpoint/%s", id.String()), url.Values{}) + if err != nil { + return GlobusEndpointInfo{}, err + } + var info GlobusEndpointInfo + err = json.Unmarshal(body, &info) + return info, err +} + +// https://docs.globus.org/api/transfer/task_submit/#get_submission_id +func (c GlobusTransferClient) getSubmissionId() (uuid.UUID, error) { + var id uuid.UUID + body, err := c.get("submission_id", url.Values{}) + if err != nil { + return id, err + } + type SubmissionIdResponse struct { + Value uuid.UUID `json:"value"` + } + var response SubmissionIdResponse + err = json.Unmarshal(body, &response) + return response.Value, err +} + +// https://docs.globus.org/api/transfer/endpoints_and_collections/#get_endpoint_or_collection_by_id +// https://docs.globus.org/api/transfer/task_submit/#submit_transfer_task +// https://docs.globus.org/api/transfer/task_submit/#transfer_item_fields +func (c GlobusTransferClient) submitTransfer(credential auth.Credential, sourceId, destinationId, submissionId uuid.UUID, + files []endpoints.FileTransfer) (uuid.UUID, error) { + var xferId uuid.UUID + + // are the source and destination endpoints configured in a conflicting way? + destinationInfo, err := c.getEndpointInfo(destinationId) + if err != nil { + return xferId, err + } + if c.Info.ForceVerify && destinationInfo.DisableVerify { // not allowed! + return xferId, &endpoints.IncompatibleDestinationError{ + Source: sourceId.String(), + SourceProvider: "globus", + Destination: destinationId.String(), + DestinationProvider: "globus", + Message: "Source endpoint forces checksum verification, but destination disables it.", + } + } + + // configure checksum settings based on destination endpoint info + var verifyChecksum bool = true + var syncLevel int = 3 // transfer only if checksums don't match + if destinationInfo.DisableVerify { // checksum verification disabled on endpoint + verifyChecksum = false + syncLevel = 2 // transfer if source file is newer than destination file + } + + type TransferItem struct { + DataType string `json:"DATA_TYPE"` // "transfer_item" + SourcePath string `json:"source_path"` + DestinationPath string `json:"destination_path"` + ExternalChecksum string `json:"external_checksum,omitempty"` + ChecksumAlgorithm string `json:"checksum_algorithm,omitempty"` + } + xferItems := make([]TransferItem, len(files)) + for i, file := range files { + var checksum, checksumAlgorithm string + if verifyChecksum { + checksum = file.Hash + checksumAlgorithm = file.HashAlgorithm + } + xferItems[i] = TransferItem{ + DataType: "transfer_item", + SourcePath: file.SourcePath, + DestinationPath: file.DestinationPath, + ExternalChecksum: checksum, + ChecksumAlgorithm: checksumAlgorithm, + } + } + + // submit the transfer request + type SubmissionRequest struct { + DataType string `json:"DATA_TYPE"` // "transfer" + Id string `json:"submission_id"` + Label string `json:"label"` // "DTS" + Data []TransferItem `json:"DATA"` + DestinationEndpoint string `json:"destination_endpoint"` + DestinationLocalUser string `json:"destination_local_user,omitempty"` + SourceEndpoint string `json:"source_endpoint"` + SyncLevel int `json:"sync_level"` + VerifyChecksum bool `json:"verify_checksum"` + FailOnQuotaErrors bool `json:"fail_on_quota_errors"` + } + data, err := json.Marshal(SubmissionRequest{ + DataType: "transfer", + Id: submissionId.String(), + Label: "DTS", + Data: xferItems, + DestinationEndpoint: destinationId.String(), + DestinationLocalUser: credential.Username, + SourceEndpoint: sourceId.String(), + SyncLevel: syncLevel, + VerifyChecksum: verifyChecksum, + FailOnQuotaErrors: true, + }) + if err != nil { + return xferId, err + } + + body, err := c.post("transfer", bytes.NewReader(data)) + if err != nil { + return xferId, err + } + type SubmissionResponse struct { + TaskId uuid.UUID `json:"task_id"` + } + + var gResp SubmissionResponse + err = json.Unmarshal(body, &gResp) + if err != nil { + return xferId, err + } + xferId = gResp.TaskId + slog.Debug(fmt.Sprintf("Initiated Globus transfer task %s (%d files)", + xferId.String(), len(files))) + return xferId, nil +} + +func (c *GlobusTransferClient) TaskStatus(taskId uuid.UUID) (GlobusTransferStatus, error) { + body, err := c.get(fmt.Sprintf("task/%s", taskId.String()), url.Values{}) + if err != nil { + return GlobusTransferStatus{}, err + } + var response GlobusTransferStatus + err = json.Unmarshal(body, &response) + return response, err +} + +func (c *GlobusTransferClient) TaskEvents(taskId uuid.UUID) ([]GlobusEvent, error) { + body, err := c.get(fmt.Sprintf("task/%s/event_list", taskId.String()), url.Values{}) + if err != nil { + return nil, err + } + type EventList struct { + Data []GlobusEvent `json:"DATA"` + } + var eventList EventList + if err = json.Unmarshal(body, &eventList); err != nil { + return nil, err + } + return eventList.Data, nil +} + +func (c *GlobusTransferClient) Cancel(taskId uuid.UUID) error { + // Because cancellation requests can't be honored under all circumstances, + // this Globus call is asynchronous. Nevertheless, the Globus documentation + // (https://docs.globus.org/api/transfer/task/#cancel_task_by_id) claims the + // call can take up to 10 seconds before returning, which doesn't meet the + // needs of the DTS. The possible outcomes of the call are identified with + // these response codes: + // 1. "Canceled", indicating that the task has been canceled + // 2. "CancelAccepted", indicating that the cancellation request has been + // acknowledged but not yet processed + // 3. "TaskComplete", indicating that the task is complete and not able to + // be canceled. + // + // We live with the 10-second wait for now, since our polling interval is + // large. + _, err := c.post(fmt.Sprintf("task/%s/cancel", taskId.String()), nil) + // NOTE: if this ^^^ becomes an issue, we can dispatch the POST to a + // NOTE: persistent goroutine to handle the cancellation + if err != nil { + if globusError, ok := err.(*GlobusTransferError); ok { + switch globusError.Code { + case "Canceled", "CancelAccepted", "TaskComplete": // it worked! + err = nil + } + } + } + return err +} + +// Uploads a file to the given (absolute) path on the HTTPS server. +func (c GlobusHttpsClient) PutFile(path string, body io.Reader) error { + resourcePath := c.Url + "/" + path + u, err := url.ParseRequestURI(resourcePath) + if err != nil { + return err + } + res := fmt.Sprintf("%v", u) + slog.Debug(fmt.Sprintf("Globus HTTPS PUT: %s", res)) + req, err := http.NewRequest(http.MethodPut, res, body) + if err != nil { + return err + } + req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + + var client http.Client + resp, err := client.Do(req) + if err != nil { + return err + } + respBody, err := io.ReadAll(resp.Body) + if err != nil { + return err + } + resp.Body.Close() + return errorFromGlobusResponse(resp, respBody) +} + +// https://docs.globus.org/globus-connect-server/v5.4/api/openapi_User_Credentials/#postUserCredential +func (c GlobusConnectServerManagerClient) AddOrUpdateUserCredential(user auth.User, provider string) (auth.Credential, error) { + for connectionProvider, credential := range user.ConnectionCredentials { + if connectionProvider == provider && provider == "s3" { + return c.addOrUpdateS3UserCredential(user, credential) + } + } + return auth.Credential{}, fmt.Errorf("unsupported user credential provider: %s", provider) +} + +//----------- +// Internals +//----------- + +// This method sends the given HTTP request, parsing the response for Globus-style error +// codes/messages and handling the ones that can be handled automatically (e.g. consent/scope +// related errors) by reauthenticating as needed and retrying the operation. See +// https://docs.globus.org/api/flows/working-with-consents for details on Globus scopes and +// consents. Returns a byte slice containing the body of the response. +func (c *GlobusTransferClient) sendRequest(request *http.Request) ([]byte, error) { + // send the initial request with a fresh HTTP client + var client http.Client + resp, err := client.Do(request) + if err != nil { + return nil, err + } + body, err := io.ReadAll(resp.Body) + if err != nil { + return nil, err + } + resp.Body.Close() + + // check the response for a Globus-style error code / message + err = errorFromGlobusResponse(resp, body) + if err != nil { + if xferErr, ok := err.(*GlobusTransferError); ok { + if xferErr.Code == "ConsentRequired" || xferErr.Code == "AuthenticationFailed" { + // our token has expired or we're missing a required scope, + // so reauthenticate + if len(xferErr.RequiredScopes) > 0 { + c.Scopes = xferErr.RequiredScopes + } + if c.AccessToken, err = c.Auth.Authenticate(c.Scopes); err != nil { + return nil, err + } + // try the request again using the new access token + request.Header.Set("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + if request.GetBody != nil { // e.g. recreate POST body + if request.Body, err = request.GetBody(); err != nil { + return nil, err + } + } + if resp, err = client.Do(request); err != nil { + return nil, err + } + body, err = io.ReadAll(resp.Body) + resp.Body.Close() + if err != nil { + return nil, err + } + return body, errorFromGlobusResponse(resp, body) + } else { + // other transfer errors are propagated + return body, err + } + } + } + return body, err +} + +// Performs a GET request on the given Globus resource, handling any obvious +// errors and returning a byte slice containing the body of the response, +// and/or any unhandled error. +func (c *GlobusTransferClient) get(resource string, values url.Values) ([]byte, error) { + resourcePath := globusTransferApiBaseUrl + fmt.Sprintf("/%s/%s", globusTransferApiVersion, resource) + u, err := url.ParseRequestURI(resourcePath) + if err != nil { + return nil, err + } + u.RawQuery = values.Encode() + res := fmt.Sprintf("%v", u) + slog.Debug(fmt.Sprintf("Globus Transfer API: GET %s", res)) + req, err := http.NewRequest(http.MethodGet, res, http.NoBody) + if err != nil { + return nil, err + } + req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + return c.sendRequest(req) +} + +// Performs a POST request on the given Globus resource, handling any obvious +// errors and returning a byte slice containing the body of the response, +// and/or any unhandled error. +// This method handles scope-related errors by reauthenticating as needed and +// retrying the operation. See https://docs.globus.org/api/flows/working-with-consents/ +// for details on Globus scopes and consents. +func (c *GlobusTransferClient) post(resource string, body io.Reader) ([]byte, error) { + resourcePath := globusTransferApiBaseUrl + fmt.Sprintf("/%s/%s", globusTransferApiVersion, resource) + u, err := url.ParseRequestURI(resourcePath) + if err != nil { + return nil, err + } + res := fmt.Sprintf("%v", u) + slog.Debug(fmt.Sprintf("Globus Transfer API: POST %s", res)) + req, err := http.NewRequest(http.MethodPost, res, body) + if err != nil { + return nil, err + } + req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + req.Header.Set("Content-Type", "application/json") + + return c.sendRequest(req) +} + +// returns an error capturing any Globus-related error in a response body, or nil if the response +// doesn't appear to be an error +func errorFromGlobusResponse(response *http.Response, body []byte) error { + bodyStr := string(body) + + // Transfer API error + if strings.Contains(bodyStr, "\"code\"") && + !strings.Contains(bodyStr, "\"code\": \"Accepted\"") && + strings.Contains(string(body), "\"message\"") { + var globusErr GlobusTransferError + err := json.Unmarshal(body, &globusErr) + if err == nil { + return &globusErr + } + } + + // Generic error + if strings.Contains(bodyStr, "GlobusError") { + return &GlobusGenericError{Message: bodyStr} + } + + // Check the status code + switch response.StatusCode { + case 200, 201: + return nil + default: + return &GlobusGenericError{Message: bodyStr} + } +} + +type GlobusEvent struct { + DataType string `json:"DATA_TYPE"` + Code string `json:"code"` + IsError bool `json:"is_error"` + Description string `json:"description"` + Details string `json:"details"` + Time string `json:"time"` +} + +func (c GlobusConnectServerManagerClient) get(resource string, values url.Values) ([]byte, error) { + resourcePath := c.Url + "/" + resource + u, err := url.ParseRequestURI(resourcePath) + if err != nil { + return nil, err + } + u.RawQuery = values.Encode() + res := fmt.Sprintf("%v", u) + slog.Debug(fmt.Sprintf("Globus Connect Server Manager API: GET %s", res)) + req, err := http.NewRequest(http.MethodGet, res, http.NoBody) + if err != nil { + return nil, err + } + req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + + var client http.Client + resp, err := client.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + return c.interpretResult(resp.Body) +} + +func (c GlobusConnectServerManagerClient) post(resource string, body io.Reader) ([]byte, error) { + resourcePath := c.Url + "/" + resource + u, err := url.ParseRequestURI(resourcePath) + if err != nil { + return nil, err + } + res := fmt.Sprintf("%v", u) + slog.Debug(fmt.Sprintf("Globus Connect Server Manager API: POST %s", res)) + req, err := http.NewRequest(http.MethodPost, res, body) + if err != nil { + return nil, err + } + req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + req.Header.Set("Content-Type", "application/json") + + var client http.Client + resp, err := client.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + return c.interpretResult(resp.Body) +} + +func (c GlobusConnectServerManagerClient) patch(resource string, body io.Reader) ([]byte, error) { + resourcePath := c.Url + "/" + resource + u, err := url.ParseRequestURI(resourcePath) + if err != nil { + return nil, err + } + res := fmt.Sprintf("%v", u) + slog.Debug(fmt.Sprintf("Globus Connect Server Manager API: PATCH %s", res)) + req, err := http.NewRequest(http.MethodPatch, res, body) + if err != nil { + return nil, err + } + req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", c.AccessToken)) + req.Header.Set("Content-Type", "application/json") + + var client http.Client + resp, err := client.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + return c.interpretResult(resp.Body) +} + +func (m GlobusConnectServerManagerClient) interpretResult(body io.Reader) (json.RawMessage, error) { + var payload []byte + var err error + if payload, err = io.ReadAll(body); err != nil { + return []byte{}, err + } + var result GlobusManagerApiResult_1_1_0 + if err = json.Unmarshal(payload, &result); err != nil { + return []byte{}, err + } + if result.HttpResponseCode != http.StatusOK && result.HttpResponseCode != http.StatusCreated { + if result.AuthorizationParameters != nil { + var params GlobusAuthorizationParameters + if err := json.Unmarshal(result.AuthorizationParameters, ¶ms); err != nil { + return []byte{}, err + } + return []byte{}, &GlobusAuthRequirementsError{ + AuthorizationParameters: params, + Message: result.Message, + Code: result.Code, + } + } + return []byte{}, errors.New(result.Message) + } + return result.Data, nil +} + +type GlobusManagerApiResult_1_1_0 struct { + DataType string `json:"DATA_TYPE"` // always `result#1.1.0` + AuthorizationParameters json.RawMessage `json:"authorization_parameters,omitempty"` // diagnostics + Code string `json:"code"` + Data json.RawMessage `json:"data"` + //Detail any `json:"detail"` + //HasNextPage bool `json:"has_next_page"` + HttpResponseCode int `json:"http_response_code"` + //Marker string `json:"marker"` + Message string `json:"message"` +} + +type GlobusS3KeysPrefixPaths_1_0_0 struct { + PathPrefixes []string `json:"path_prefixes"` + S3KeyId string `json:"s3_key_id"` + S3SecretKey string `json:"s3_secret_key"` +} +type GlobusS3UserCredentialPolicies_1_2_0 struct { + DataType string `json:"DATA_TYPE"` // always `s3_user_credential_policies#1.2.0` + S3KeyId string `json:"s3_key_id"` + S3MultiKeys []GlobusS3KeysPrefixPaths_1_0_0 `json:"s3_multi_keys,omitempty"` + S3RequesterPays bool `json:"s3_requester_pays,omitempty"` + S3SecretKey string `json:"s3_secret_key"` +} +type GlobusUserCredentialRecord struct { + DataType string `json:"DATA_TYPE"` // always `user_credential#1.0.0` + ConnectorId string `json:"connector_id,omitempty"` + Deleted bool `json:"deleted"` + DisplayName string `json:"display_name,omitempty"` + Id string `json:"id"` + IdentityId string `json:"identity_id"` + Invalid bool `json:"invalid"` + Policies json.RawMessage `json:"policies"` + Provisioned bool `json:"provisioned,omitempty"` + StorageGatewayId string `json:"storage_gateway_id"` + Username string `json:"username"` +} + +func (m *GlobusConnectServerManagerClient) getStorageGatewayInfo() error { + data, err := m.get("api/storage_gateways/", url.Values{}) + if err != nil { + return err + } + type GlobusStorageGateway_1_3_0 struct { + ConnectorId string `json:"connector_id"` + DataType string `json:"DATA_TYPE"` // always `s3_storage_gateway#1.3.0` + Id string `json:"id"` + Policies json.RawMessage `json:"policies"` + } + var gateways []GlobusStorageGateway_1_3_0 + if err := json.Unmarshal(data, &gateways); err != nil { + return err + } + for _, g := range gateways { + var gateway GlobusStorageGateway + gateway.ConnectorId = uuid.MustParse(g.ConnectorId) + gateway.Id = uuid.MustParse(g.Id) + /* + type GlobusS3StoragePolicies_1_3_0 struct { + DataType string `json:"DATA_TYPE"` // always `s3_storage_policies#1.3.0` + S3Buckets string `json:"s3_buckets"` + S3Endpoint string `json:"s3_endpoint"` + } + var policy GlobusS3StoragePolicies_1_3_0 + if err = json.Unmarshal(g.Policies, &policy); err != nil { + continue + } + */ + gateway.Provider = "s3" + m.StorageGateways = append(m.StorageGateways, gateway) + } + return nil +} + +// NOTE: For now, we only allow a single S3 credential per user to be registered with a Globus +// NOTE: endpoint per user, using the user's ORCID. +func (m *GlobusConnectServerManagerClient) addOrUpdateS3UserCredential(user auth.User, s3Credential auth.Credential) (auth.Credential, error) { + var record GlobusUserCredentialRecord + var found bool + var payload []byte + var err error + + globusCred, foundGlobusId := user.ConnectionCredentials["globus"] + if !foundGlobusId { + return auth.Credential{}, fmt.Errorf("no Globus ID is associated with this user") + } + slog.Debug(fmt.Sprintf("Globus ID: %s", globusCred.Id)) + slog.Debug(fmt.Sprintf("Globus Username: %s", globusCred.Username)) + + if record, found, _ = m.findUserCredentialRecord(globusCred); found { + // Update the record with an S3 policy + slog.Debug("Looking for user S3 credential...") + var s3Policy GlobusS3UserCredentialPolicies_1_2_0 + err := json.Unmarshal(record.Policies, &s3Policy) + if err != nil || s3Policy.S3KeyId != s3Credential.Id || s3Policy.S3SecretKey != s3Credential.Secret { + // insert an S3 policy and patch the registered credential + s3Policy.DataType = "s3_user_credential_policies#1.2.0" + s3Policy.S3KeyId = s3Credential.Id + s3Policy.S3SecretKey = s3Credential.Secret + if record.Policies, err = json.Marshal(s3Policy); err != nil { + return auth.Credential{}, err + } + if payload, err = json.Marshal(record); err != nil { + return auth.Credential{}, err + } + if _, err = m.patch("api/user_credentials", bytes.NewReader(payload)); err != nil { + return auth.Credential{}, err + } + } + } else { + // No existing record -- create a new one. + var newS3Policy []byte + if newS3Policy, err = json.Marshal(GlobusS3UserCredentialPolicies_1_2_0{ + DataType: "s3_user_credential_policies#1.2.0", + S3KeyId: s3Credential.Id, + S3SecretKey: s3Credential.Secret, + }); err != nil { + return auth.Credential{}, err + } + + // Attempt to register the S3 credential with our storage gateways until one accepts. + // FIXME: This is where the remaining auth issue is. The error message I encounter with + // FIXME: the correct gateway is: `Identity set contains an identity from an allowed domain, + // FIXME: but it does not map to a valid username for this connector`. This suggests to me that + // FIXME: either the user's Globus ID (globusCred.Id) or the mapped username (globusCred.Username) + // FIXME: is incorrect, but I've checked my own account's values against the ORCID identity + // FIXME: shown at https://app.globus.org/settings/identities and they are correct. + registrations := 0 + for _, gateway := range m.StorageGateways { + if gateway.Provider == "s3" { + record = GlobusUserCredentialRecord{ + DataType: "user_credential#1.0.0", + //ConnectorId: gateway.ConnectorId.String(), + //DisplayName: user.Name, + Id: globusCred.Id, + IdentityId: m.ClientId, + Policies: newS3Policy, + //Provisioned: true, + StorageGatewayId: gateway.Id.String(), + Username: globusCred.Username, + } + if payload, err = json.Marshal(record); err != nil { + return auth.Credential{}, err + } + _, err = m.post("api/user_credentials", bytes.NewReader(payload)) + if err != nil { + slog.Debug(fmt.Sprintf("Couldn't register S3 credential at storage gateway %s: %s", gateway.Id.String(), err.Error())) + } else { + // Now that we know this gateway works, eliminate the others. + m.StorageGateways = []GlobusStorageGateway{gateway} + registrations += 1 + break + } + } + } + if registrations == 0 { + return auth.Credential{}, errors.New("couldn't register an S3 credential at any storage gateway") + } + } + return s3Credential, nil +} + +func (m GlobusConnectServerManagerClient) findUserCredentialRecord(credential auth.Credential) (GlobusUserCredentialRecord, bool, error) { + for _, gateway := range m.StorageGateways { + var response GlobusManagerApiResult_1_1_0 + values := url.Values{} + values.Add("include", "all") + values.Add("storage_gateway", gateway.Id.String()) + body, err := m.get(fmt.Sprintf("api/user_credentials/%s", credential.Id), url.Values{}) + if err != nil { + return GlobusUserCredentialRecord{}, false, err + } + if err := json.Unmarshal(body, &response); err != nil { + return GlobusUserCredentialRecord{}, false, err + } + if response.HttpResponseCode != http.StatusOK && response.HttpResponseCode != http.StatusCreated { + return GlobusUserCredentialRecord{}, false, errors.New(response.Message) + } + var existingCred GlobusUserCredentialRecord + if err := json.Unmarshal(response.Data, &existingCred); err == nil { + return existingCred, true, nil + } + } + return GlobusUserCredentialRecord{}, false, nil +} + +const ( + globusTransferApiBaseUrl = "https://transfer.api.globusonline.org" + globusTransferApiVersion = "v0.10" +) diff --git a/endpoints/local/endpoint.go b/endpoints/local/endpoint.go index 13969cf0..0ec23eec 100644 --- a/endpoints/local/endpoint.go +++ b/endpoints/local/endpoint.go @@ -32,7 +32,9 @@ import ( "github.com/google/uuid" "github.com/mitchellh/mapstructure" + "github.com/kbase/dts/auth" "github.com/kbase/dts/endpoints" + "github.com/kbase/dts/endpoints/globus" "github.com/kbase/dts/endpoints/s3" ) @@ -48,26 +50,26 @@ type Endpoint struct { // descriptive endpoint name (obtained from config) Name string // endpoint UUID (obtained from config) - Id uuid.UUID - // root directory for endpoint (default: current working directory) - root string + Id_ uuid.UUID + Paths struct { + Base string + Data string + } // transfers in progress Xfers map[uuid.UUID]xferRecord } // configuration struct for local endpoint type Config struct { - Name string `yaml:"name"` - Id string `yaml:"id"` - Root string `yaml:"root"` + Name string `yaml:"name"` + Id string `yaml:"id"` + BasePath string `yaml:"base_path" mapstructure:"base_path,omitempty"` + DataPath string `yaml:"data_path" mapstructure:"data_path,omitempty"` } // creates a new local endpoint using the information supplied in the // DTS configuration file under the given endpoint name func NewEndpoint(config Config) (endpoints.Endpoint, error) { - if config.Root == "" { - config.Root = "/" - } if config.Name == "" { return nil, fmt.Errorf("name must be specified for local endpoint") } @@ -77,10 +79,10 @@ func NewEndpoint(config Config) (endpoints.Endpoint, error) { } ep := &Endpoint{ Name: config.Name, - Id: id, + Id_: id, Xfers: make(map[uuid.UUID]xferRecord), } - err = ep.setRoot(config.Root) + err = ep.setPaths(config.BasePath, config.DataPath) return ep, err } @@ -94,25 +96,54 @@ func EndpointConstructor(conf map[string]any) (endpoints.Endpoint, error) { } // sets the root directory for the local endpoint after checking that it exists -func (ep *Endpoint) setRoot(dir string) error { - _, err := os.Stat(dir) - if err == nil { - ep.root = dir +func (ep *Endpoint) setPaths(base, data string) error { + if base == "" { + ep.Paths.Base = "/" + } else { + _, err := os.Stat(base) + if err != nil { + return fmt.Errorf("couldn't set base path '%s' for local endpoint: %s", base, err.Error()) + } + ep.Paths.Base = base + } + if data != "" { + dataPath := filepath.Join(ep.Paths.Base, data) + if _, err := os.Stat(dataPath); err != nil { + return fmt.Errorf("couldn't set data path '%s' for local endpoint: %s", dataPath, err.Error()) + } } - return err + ep.Paths.Data = data + return nil } -func (ep *Endpoint) Provider() string { +func (ep Endpoint) Id() uuid.UUID { + return ep.Id_ +} + +func (ep Endpoint) Provider() string { return "local" } -func (ep *Endpoint) Root() string { - return ep.root +func (ep Endpoint) BasePath() string { + return ep.Paths.Base +} + +func (ep Endpoint) DataPath() string { + return ep.Paths.Data } -func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { +func (ep Endpoint) ConnectsWith(provider string) bool { + switch provider { + case "s3", "globus": + return true + default: + return false + } +} + +func (ep Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { for _, descriptor := range descriptors { - absPath := filepath.Join(ep.root, descriptor["path"].(string)) + absPath := filepath.Join(ep.BasePath(), ep.DataPath(), descriptor["path"].(string)) _, err := os.Stat(absPath) if err != nil { return false, nil @@ -121,7 +152,7 @@ func (ep *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { return true, nil } -func (ep *Endpoint) Transfers() ([]uuid.UUID, error) { +func (ep Endpoint) Transfers() ([]uuid.UUID, error) { xfers := make([]uuid.UUID, 0) for xferId, xfer := range ep.Xfers { switch xfer.Status.Code { @@ -142,41 +173,7 @@ func (ep *Endpoint) transferFiles(xferId uuid.UUID, dest endpoints.Endpoint) { if xfer.Canceled { break } - - sourcePath := filepath.Join(ep.Root(), file.SourcePath) - destPath := filepath.Join(dest.Root(), file.DestinationPath) - - // check for the source directory - sourceDir := filepath.Dir(sourcePath) - var sourceDirInfo os.FileInfo - sourceDirInfo, err = os.Stat(sourceDir) - if err != nil { - break - } - - // create the destination directory if needed - destDir := filepath.Dir(destPath) - _, err = os.Stat(destDir) - if err != nil { - if errors.Is(err, fs.ErrNotExist) { // destination dir doesn't exist - os.MkdirAll(destDir, sourceDirInfo.Mode()) - } else { // something else happened - break - } - } - - // copy the file into place - var data []byte - var sourceFileInfo os.FileInfo - sourceFileInfo, err = os.Stat(sourcePath) - if err != nil { - break - } - data, err = os.ReadFile(sourcePath) - if err != nil { - break - } - err = os.WriteFile(destPath, data, sourceFileInfo.Mode()) + err = ep.transferFile(dest, file) if err != nil { break } @@ -193,12 +190,50 @@ func (ep *Endpoint) transferFiles(xferId uuid.UUID, dest endpoints.Endpoint) { ep.Xfers[xferId] = xfer } -func (ep *Endpoint) Transfer(dst endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { +// implements per-file local transfers and validation +func (ep *Endpoint) transferFile(dest endpoints.Endpoint, file endpoints.FileTransfer) error { + sourcePath := filepath.Join(ep.BasePath(), ep.DataPath(), file.SourcePath) + destPath := filepath.Join(dest.BasePath(), dest.DataPath(), file.DestinationPath) + + // check for the source directory + sourceDir := filepath.Dir(sourcePath) + sourceDirInfo, err := os.Stat(sourceDir) + if err != nil { + return err + } + + // create the destination directory if needed + destDir := filepath.Dir(destPath) + _, err = os.Stat(destDir) + if err != nil { + if errors.Is(err, fs.ErrNotExist) { // destination dir doesn't exist + os.MkdirAll(destDir, sourceDirInfo.Mode()) + } else { // something else happened + return err + } + } + + // copy the file into place + var data []byte + var sourceFileInfo os.FileInfo + sourceFileInfo, err = os.Stat(sourcePath) + if err != nil { + return err + } + data, err = os.ReadFile(sourcePath) + if err != nil { + return err + } + return os.WriteFile(destPath, data, sourceFileInfo.Mode()) +} + +func (ep *Endpoint) Transfer(user auth.User, dst endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { var xferId uuid.UUID _, isLocal := dst.(*Endpoint) _, isS3 := dst.(*s3.Endpoint) - if !isLocal && !isS3 { + _, isGlobus := dst.(*globus.Endpoint) + if !isLocal && !isS3 && !isGlobus { return xferId, &endpoints.IncompatibleDestinationError{ Source: ep.Name, SourceProvider: "local", @@ -223,21 +258,23 @@ func (ep *Endpoint) Transfer(dst endpoints.Endpoint, files []endpoints.FileTrans } // all files are staged; start the transfer - if isS3 { - // special case: destination is S3 endpoint - // turn each file into a bytes.Reader and upload it + if isS3 || isGlobus { + // upload each file via PUT for _, file := range files { - sourcePath := filepath.Join(ep.Root(), file.SourcePath) + sourcePath := filepath.Join(ep.BasePath(), ep.DataPath(), file.SourcePath) data, err := os.ReadFile(sourcePath) if err != nil { - err = fmt.Errorf("incomplete file transfer at: %s for S3 transfer: %w", sourcePath, err) + err = fmt.Errorf("incomplete file transfer: couldn't transfer %s to %s endpoint: %w", sourcePath, dst.Provider(), err) return xferId, err } reader := bytes.NewReader(data) - s3Dst := dst.(*s3.Endpoint) - err = s3Dst.PutFromReader(file.DestinationPath, reader) + if s3Dst, ok := dst.(*s3.Endpoint); ok { + err = s3Dst.PutFromReader(file.DestinationPath, reader) + } else if globusDst, ok := dst.(*globus.Endpoint); ok { + err = globusDst.PutFromReader(file.DestinationPath, reader) + } if err != nil { - err = fmt.Errorf("incomplete file transfer at: %s for S3 transfer: %w", file.DestinationPath, err) + err = fmt.Errorf("incomplete file transfer: couldn't transfer %s to %s endpoint: %w", file.DestinationPath, dst.Provider(), err) return xferId, err } } @@ -254,7 +291,7 @@ func (ep *Endpoint) Transfer(dst endpoints.Endpoint, files []endpoints.FileTrans return xferId, nil } - // non-S3 endpoints are handled entirely within local endpoint + // non-S3/Globus endpoints are handled entirely within local endpoint // assign a UUID to the transfer and set it going xferId = uuid.New() ep.Xfers[xferId] = xferRecord{ @@ -267,7 +304,6 @@ func (ep *Endpoint) Transfer(dst endpoints.Endpoint, files []endpoints.FileTrans } go ep.transferFiles(xferId, dst) return xferId, nil - } func (ep *Endpoint) Status(id uuid.UUID) (endpoints.TransferStatus, error) { @@ -290,5 +326,5 @@ func (ep *Endpoint) Cancel(id uuid.UUID) error { // this method is specific to local endpoints and gives access to the // local filesystem func (ep *Endpoint) FS() (fs.FS, error) { - return os.DirFS(filepath.Join("/", ep.root)), nil + return os.DirFS(filepath.Join(ep.BasePath(), ep.DataPath())), nil } diff --git a/endpoints/local/endpoint_test.go b/endpoints/local/endpoint_test.go index 78453edd..0cd0801e 100644 --- a/endpoints/local/endpoint_test.go +++ b/endpoints/local/endpoint_test.go @@ -31,6 +31,7 @@ import ( "github.com/stretchr/testify/assert" "github.com/stretchr/testify/assert/yaml" + "github.com/kbase/dts/auth" "github.com/kbase/dts/endpoints" ) @@ -87,19 +88,19 @@ func setup() { name: Source Endpoint id: 2ee69538-10d5-4d1e-a890-1127b5e42003 provider: local -root: %s +base_path: %s `, sourceRoot) destConfig = fmt.Sprintf(` name: Destination Endpoint id: b925d96e-7e39-473b-a658-714f8c243b1c provider: local -root: %s +base_path: %s `, destinationRoot) destCancelConfig = fmt.Sprintf(` name: Destination Endpoint for cancellation id: b925d96e-7e39-473b-a658-714f8c243b1c provider: local -root: %s +base_path: %s `, destinationRootCancel) } @@ -124,9 +125,9 @@ func TestBadLocalConstructor(t *testing.T) { assert := assert.New(t) conf := Config{ - Name: "", - Id: uuid.New().String(), - Root: "/bad/endpoint/no/name", + Name: "", + Id: uuid.New().String(), + BasePath: "/bad/endpoint/no/name", } endpoint, err := NewEndpoint(conf) assert.Nil(endpoint) @@ -205,7 +206,7 @@ func TestLocalTransfer(t *testing.T) { DestinationPath: sourceFilesById[id], }) } - _, err = source.Transfer(destination, fileXfers) + _, err = source.Transfer(auth.User{}, destination, fileXfers) assert.Nil(err) } @@ -229,7 +230,7 @@ func TestBadLocalTransfer(t *testing.T) { DestinationPath: sourceFilesById[id] + "_with_bad_suffix", }) } - _, err = source.Transfer(destination, fileXfers) + _, err = source.Transfer(auth.User{}, destination, fileXfers) assert.NotNil(err) } @@ -268,7 +269,7 @@ func TestLocalTransferCancellation(t *testing.T) { DestinationPath: sourceFilesById[id], }) } - id, err := source.Transfer(destination, fileXfers) + id, err := source.Transfer(auth.User{}, destination, fileXfers) assert.Nil(err) err = source.Cancel(id) assert.Nil(err) diff --git a/endpoints/s3/endpoint.go b/endpoints/s3/endpoint.go index dc5967b2..9d98bc6e 100644 --- a/endpoints/s3/endpoint.go +++ b/endpoints/s3/endpoint.go @@ -37,6 +37,7 @@ import ( "github.com/google/uuid" "github.com/mitchellh/mapstructure" + "github.com/kbase/dts/auth" "github.com/kbase/dts/endpoints" ) @@ -70,7 +71,7 @@ type Endpoint struct { // AWS S3 uploader Uploader *manager.Uploader // endpoint UUID (obtained from config) - Id uuid.UUID + Id_ uuid.UUID // Map of completed transfers TransfersMap map[uuid.UUID]*TransferStatus } @@ -125,7 +126,7 @@ func NewEndpoint(bucket string, id uuid.UUID, ecfg Config) (endpoints.Endpoint, newEndpoint.Downloader = manager.NewDownloader(newEndpoint.Client) newEndpoint.Uploader = manager.NewUploader(newEndpoint.Client) newEndpoint.Bucket = bucket - newEndpoint.Id = id + newEndpoint.Id_ = id newEndpoint.TransfersMap = make(map[uuid.UUID]*TransferStatus) return &newEndpoint, nil @@ -148,14 +149,27 @@ func EndpointConstructor(conf map[string]any) (endpoints.Endpoint, error) { return NewEndpoint(config.Bucket, id, config.Config) } -func (e *Endpoint) Provider() string { +func (e Endpoint) Id() uuid.UUID { + return e.Id_ +} + +func (e Endpoint) Provider() string { return "s3" } -func (e *Endpoint) Root() string { +func (e Endpoint) BasePath() string { return e.Bucket + "/" } +func (e *Endpoint) DataPath() string { + return "" +} + +func (e *Endpoint) ConnectsWith(provider string) bool { + // The S3 endpoint can't send to anyone else at the moment. + return provider == "s3" +} + func (e *Endpoint) FilesStaged(descriptors []map[string]any) (bool, error) { staged := true for _, d := range descriptors { @@ -187,7 +201,7 @@ func (e *Endpoint) Transfers() ([]uuid.UUID, error) { return ids, nil } -func (e *Endpoint) Transfer(dst endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { +func (e *Endpoint) Transfer(user auth.User, dst endpoints.Endpoint, files []endpoints.FileTransfer) (uuid.UUID, error) { s3Dest, ok := dst.(*Endpoint) if !ok { return uuid.Nil, fmt.Errorf("destination endpoint is not an S3 endpoint") diff --git a/endpoints/s3/endpoint_test.go b/endpoints/s3/endpoint_test.go index 14c1b040..3e735d42 100644 --- a/endpoints/s3/endpoint_test.go +++ b/endpoints/s3/endpoint_test.go @@ -36,6 +36,7 @@ import ( "github.com/google/uuid" "github.com/stretchr/testify/assert" + "github.com/kbase/dts/auth" "github.com/kbase/dts/endpoints" ) @@ -145,7 +146,7 @@ func TestNewAWSS3Endpoint(t *testing.T) { awsEndpoint, err := NewEndpoint(awsTestBucket, uuid.New(), cfg) assert.NotNil(awsEndpoint) assert.Nil(err) - assert.Equal(awsTestBucket+"/", awsEndpoint.Root()) + assert.Equal(awsTestBucket+"/", awsEndpoint.BasePath()) assert.Equal("s3", awsEndpoint.Provider()) staged, err := awsEndpoint.FilesStaged([]map[string]any{}) assert.True(staged) @@ -179,7 +180,7 @@ func TestNewMinioS3Endpoint(t *testing.T) { minioEndpoint, err := NewEndpoint(minioTestBuckets[0], uuid.New(), cfg) assert.NotNil(minioEndpoint) assert.Nil(err) - assert.Equal(minioTestBuckets[0]+"/", minioEndpoint.Root()) + assert.Equal(minioTestBuckets[0]+"/", minioEndpoint.BasePath()) assert.Equal("s3", minioEndpoint.Provider()) // test FilesStaged with existing files @@ -275,7 +276,7 @@ func TestAWSToMinioTransfer(t *testing.T) { DestinationPath: "LICENSE_copied.txt", }, } - transferID, err := awsEndpoint.Transfer(minioEndpoint, filesToTransfer) + transferID, err := awsEndpoint.Transfer(auth.User{}, minioEndpoint, filesToTransfer) assert.NotEqual(uuid.Nil, transferID) assert.Nil(err) @@ -349,7 +350,7 @@ func TestMinioToMinioTransfer(t *testing.T) { DestinationPath: "testfile2_copied.txt", }, } - transferID, err := minioSrcEndpoint.Transfer(minioDestEndpoint, filesToTransfer) + transferID, err := minioSrcEndpoint.Transfer(auth.User{}, minioDestEndpoint, filesToTransfer) assert.NotEqual(uuid.Nil, transferID) assert.Nil(err) @@ -407,7 +408,7 @@ func TestMinioToMinioTransfer(t *testing.T) { DestinationPath: "testfile1_copied_again.txt", }, } - failedTransferID, err := minioSrcEndpoint.Transfer(minioDestEndpoint, nonexistentFileTransfer) + failedTransferID, err := minioSrcEndpoint.Transfer(auth.User{}, minioDestEndpoint, nonexistentFileTransfer) assert.NotEqual(uuid.Nil, failedTransferID) assert.Nil(err) @@ -470,7 +471,7 @@ func TestMinioToMinioTransfer(t *testing.T) { DestinationPath: "testfile3_copied.txt", }, } - cancelTransferID, err := minioSrcEndpoint.Transfer(minioDestEndpoint, allFilesTransfer) + cancelTransferID, err := minioSrcEndpoint.Transfer(auth.User{}, minioDestEndpoint, allFilesTransfer) assert.NotEqual(uuid.Nil, cancelTransferID) assert.Nil(err) diff --git a/integration/irods/fixtures/test-config.yaml b/integration/irods/fixtures/test-config.yaml index b17d5a0f..31a52c0c 100644 --- a/integration/irods/fixtures/test-config.yaml +++ b/integration/irods/fixtures/test-config.yaml @@ -15,7 +15,7 @@ endpoints: name: local-fs id: 550e8400-e29b-41d4-a716-446655440000 provider: local - root: . + base_path: . s3-foo: id: 6ba7b810-9dad-11d1-80b4-00c04fd430c8 bucket: test-bucket-integration-irods-foo diff --git a/integration/minio/fixtures/test-config.yaml b/integration/minio/fixtures/test-config.yaml index d1e7f23a..4710ea1b 100644 --- a/integration/minio/fixtures/test-config.yaml +++ b/integration/minio/fixtures/test-config.yaml @@ -15,7 +15,7 @@ endpoints: name: local-fs id: 550e8400-e29b-41d4-a716-446655440000 provider: local - root: . + base_path: . s3-foo: id: 6ba7b810-9dad-11d1-80b4-00c04fd430c8 bucket: test-bucket-integration-foo diff --git a/services/prototype.go b/services/prototype.go index 7b8eb5ef..4483f3e7 100644 --- a/services/prototype.go +++ b/services/prototype.go @@ -160,10 +160,10 @@ func authorize(authorizationHeader string) (auth.User, error) { } } if err != nil { + // maybe it's a KBase token, so check with the KBase auth server slog.Debug(fmt.Sprintf("authenticator: %s", err.Error())) slog.Debug("Falling back to KBase authentication.") - // maybe it's a KBase dev token, so check with the KBase auth server authServer, err := auth.NewKBaseAuthServer(accessToken) if err != nil { return auth.User{}, huma.Error401Unauthorized(err.Error()) diff --git a/services/prototype_test.go b/services/prototype_test.go index a44a7c6d..6504061e 100644 --- a/services/prototype_test.go +++ b/services/prototype_test.go @@ -96,17 +96,17 @@ endpoints: name: Endpoint 1 id: 26d61236-39f6-4742-a374-8ec709347f2f provider: local - root: SOURCE_ROOT + base_path: SOURCE_ROOT destination-endpoint1: name: Endpoint 2 id: f1865b86-2c64-4b8b-99f3-5aaa945ec3d9 provider: local - root: DESTINATION1_ROOT + base_path: DESTINATION1_ROOT destination-endpoint2: name: Endpoint 3 id: f1865b86-2c64-4b8b-99f3-5aaa945ec3d9 provider: local - root: DESTINATION2_ROOT + base_path: DESTINATION2_ROOT ` // file test metadata diff --git a/services/version.go b/services/version.go index 85383a01..e1a5f049 100644 --- a/services/version.go +++ b/services/version.go @@ -6,8 +6,8 @@ import ( // Version numbers var majorVersion = 0 -var minorVersion = 13 -var patchVersion = 4 +var minorVersion = 15 +var patchVersion = 0 // Version string var version = fmt.Sprintf("%d.%d.%d", majorVersion, minorVersion, patchVersion) diff --git a/transfers/dispatcher.go b/transfers/dispatcher.go index 1a221182..1a49621d 100644 --- a/transfers/dispatcher.go +++ b/transfers/dispatcher.go @@ -268,7 +268,7 @@ func (d *dispatcherState) initialize(transferId uuid.UUID) error { return NoFilesAvailableError{Endpoint: spec.Source} } - // do we need to stage files for the source database? + // Do we need to stage files for the source database? filesStaged := true descriptorsForEndpoint, err := descriptorsByEndpoint(spec, descriptors) if err != nil { @@ -288,6 +288,7 @@ func (d *dispatcherState) initialize(transferId uuid.UUID) error { } } + // Get moving. if !filesStaged { err = stager.StageFiles(transferId) } else { diff --git a/transfers/manifestor.go b/transfers/manifestor.go index 6be84d1b..91343d08 100644 --- a/transfers/manifestor.go +++ b/transfers/manifestor.go @@ -251,7 +251,7 @@ func (m *manifestorState) generateAndSendManifest(transferId uuid.UUID) (manifes if err != nil { return manifestEntry{}, err } - manifestXferId, err := source.Transfer(destination, []FileTransfer{ + manifestXferId, err := source.Transfer(spec.User, destination, []FileTransfer{ { SourcePath: filename, DestinationPath: filepath.Join(destinationFolder, "manifest.json"), diff --git a/transfers/mover.go b/transfers/mover.go index 2fe9fb09..267db057 100644 --- a/transfers/mover.go +++ b/transfers/mover.go @@ -275,7 +275,22 @@ func (m *moverState) moveFiles(transferId uuid.UUID) ([]moveOperation, error) { if err != nil { return nil, err } - moveId, err := sourceEndpoint.Transfer(destinationEp, files) + + // Handle connections between endpoints with different providers. + if sourceEndpoint.Provider() != destinationEp.Provider() { + if !sourceEndpoint.ConnectsWith(destinationEp.Provider()) { + return nil, &endpoints.IncompatibleDestinationError{ + Source: source, + SourceProvider: sourceEndpoint.Provider(), + Destination: spec.Destination, + DestinationProvider: destinationEp.Provider(), + Message: fmt.Sprintf("a %s endpoints cannot transfer files to a %s endpoint", + sourceEndpoint.Provider(), destinationEp.Provider()), + } + } + } + + moveId, err := sourceEndpoint.Transfer(spec.User, destinationEp, files) if err != nil { return nil, err } diff --git a/transfers/store.go b/transfers/store.go index 1d7be4e7..005c5f3a 100644 --- a/transfers/store.go +++ b/transfers/store.go @@ -25,13 +25,16 @@ import ( "cmp" "encoding/gob" "fmt" + "log/slog" "slices" "time" "github.com/google/uuid" + "github.com/kbase/dts/auth" "github.com/kbase/dts/config" "github.com/kbase/dts/databases" + "github.com/kbase/dts/databases/kbase_lakehouse" // for Globus S3 connector HACK ) //------- @@ -248,9 +251,8 @@ func (s *storeState) process(decoder *gob.Decoder) { Time: time.Now(), }) } else { - size := transfers[id].payloadSize() publish(Message{ - Description: fmt.Sprintf("Created new transfer %s (%d file(s), %g GB)", id, newXfer.Status.NumFiles, float64(size)/float64(1024*1024*1024)), + Description: fmt.Sprintf("Created new transfer %s (%d file(s))", id, newXfer.Status.NumFiles), TransferId: id, TransferStatus: transfers[id].Status, Time: time.Now(), @@ -316,6 +318,7 @@ func (s *storeState) process(decoder *gob.Decoder) { } } case encoder := <-s.Channels.SaveAndStop: + s.eraseConnectionCredentials(transfers) s.Channels.Error <- encoder.Encode(transfers) running = false } @@ -382,6 +385,61 @@ func (s *storeState) newTransfer(spec Specification) transferStoreEntry { slices.SortFunc(descriptors, func(a, b map[string]any) int { return cmp.Compare(a["id"].(string), b["id"].(string)) }) + + // Determine all source endpoints. + sourceEndpoints := make(map[string]bool) + for _, d := range descriptors { + var endpointName string + entry, keyFound := d["endpoint"] + if keyFound { + endpointName, _ = entry.(string) + } + if endpointName == "" { + endpointName = source.EndpointNames()[0] + } + if _, endpointFound := sourceEndpoints[endpointName]; !endpointFound { + sourceEndpoints[endpointName] = true + } + } + + // HACK: Special logic for Globus transfers to KBase Lakehouse via S3 Connector: + // HACK: A Globus ID іs required for every user for which we register S3 credentials for + // HACK: connectors. We attempt to fetch this ID from the KBase Lakehouse database + { + dest, err := databases.NewDatabase(spec.Destination) + if err != nil { + return transferStoreEntry{ + Spec: spec, + Status: TransferStatus{ + Code: TransferStatusFailed, + Message: err.Error(), + NumFiles: len(spec.FileIds), + }, + } + } + if kbLakehouse, ok := dest.(*kbase_lakehouse.Database); ok { + slog.Debug("Extracting Globus ID for user") + globusId, err := kbLakehouse.GlobusId(spec.User.Orcid) + if err != nil { + return transferStoreEntry{ + Spec: spec, + Status: TransferStatus{ + Code: TransferStatusFailed, + Message: err.Error(), + NumFiles: len(spec.FileIds), + }, + } + } + if globusId.String() != "" { + slog.Debug(fmt.Sprintf("Adding Globus ID %s for user", globusId.String())) + spec.User.ConnectionCredentials["globus"] = auth.Credential{ + Id: globusId.String(), + Username: fmt.Sprintf("%s@orcid.org", spec.User.Orcid), + } + } + } + } + entry := transferStoreEntry{ Descriptors: descriptors, Spec: spec, @@ -392,3 +450,11 @@ func (s *storeState) newTransfer(spec Specification) transferStoreEntry { return entry } + +// clears user connection credentials from transfer specifications so they don't get written to disk +func (s *storeState) eraseConnectionCredentials(transfers map[uuid.UUID]transferStoreEntry) { + for i, transfer := range transfers { + transfer.Spec.User.ConnectionCredentials = map[string]auth.Credential{} + transfers[i] = transfer + } +} diff --git a/transfers/transfers.go b/transfers/transfers.go index cd9c63fc..08c7e8d9 100644 --- a/transfers/transfers.go +++ b/transfers/transfers.go @@ -39,6 +39,7 @@ import ( "github.com/kbase/dts/databases" "github.com/kbase/dts/databases/jdp" "github.com/kbase/dts/databases/kbase" + "github.com/kbase/dts/databases/kbase_lakehouse" "github.com/kbase/dts/databases/nmdc" s3db "github.com/kbase/dts/databases/s3" "github.com/kbase/dts/endpoints" @@ -225,12 +226,12 @@ type resultType[V any] struct { func registerEndpointProviders() error { // NOTE: it's okay if these endpoint providers have already been registered, // NOTE: as they can be used in testing - endpointsToRegister := map[string]func(conf map[string]any) (endpoints.Endpoint, error){ + providersToRegister := map[string]func(conf map[string]any) (endpoints.Endpoint, error){ "globus": globus.EndpointConstructor, "local": local.EndpointConstructor, "s3": s3ep.EndpointConstructor, } - for name, constructor := range endpointsToRegister { + for name, constructor := range providersToRegister { err := endpoints.RegisterEndpointProvider(name, constructor) if err != nil { // ignore AlreadyRegisteredError but propagate others @@ -242,13 +243,6 @@ func registerEndpointProviders() error { return nil } -// constructors for named (bespoke) databases -var dbConstructors map[string]func(config map[string]any) func() (databases.Database, error) = map[string]func(config map[string]any) func() (databases.Database, error){ - "jdp": jdp.DatabaseConstructor, - "kbase": kbase.DatabaseConstructor, - "nmdc": nmdc.DatabaseConstructor, -} - // registers databases; if at least one database is available, no error is propagated func registerDatabases(conf config.Config) error { for dbName, dbConf := range conf.Databases { @@ -279,6 +273,13 @@ func registerDatabases(conf config.Config) error { slog.Debug(fmt.Sprintf("No 'delete_after' pruning time specified for database '%s'; using default of %d", dbName, conf.Service.DeleteAfter)) dbConf["delete_after"] = conf.Service.DeleteAfter } + dbConstructors := map[string]func(config map[string]any) func() (databases.Database, error){ + "jdp": jdp.DatabaseConstructor, + "kbase": kbase.DatabaseConstructor, + "kbase_lakehouse": kbase_lakehouse.DatabaseConstructor, + "nmdc": nmdc.DatabaseConstructor, + } + if constructor, found := dbConstructors[dbName]; found { if err := databases.RegisterDatabase(dbName, constructor(dbConf)); err != nil { slog.Error(err.Error()) @@ -417,9 +418,9 @@ func determineDestinationEndpoint(destination string) (endpoints.Endpoint, error return nil, err } conf := globus.Config{ - Name: fmt.Sprintf("Custom endpoint (%s)", endpointId.String()), - Id: endpointId.String(), - Root: customSpec.Path, + Name: fmt.Sprintf("Custom endpoint (%s)", endpointId.String()), + Id: endpointId.String(), + DataPath: customSpec.Path, Credential: auth.Credential{ Id: clientId.String(), Secret: credential.Secret, diff --git a/transfers/transfers_test.go b/transfers/transfers_test.go index 7fee586c..90d8f796 100644 --- a/transfers/transfers_test.go +++ b/transfers/transfers_test.go @@ -230,12 +230,10 @@ endpoints: name: Endpoint 1 id: 26d61236-39f6-4742-a374-8ec709347f2f provider: test - root: SOURCE_ROOT destination-endpoint: name: Endpoint 2 id: f1865b86-2c64-4b8b-99f3-5aaa945ec3d9 provider: test - root: DESTINATION_ROOT ` var testDescriptors map[string]map[string]any = map[string]map[string]any{