Updated logic to handle Snapshot based bootstrapping. The apidconfig will be tracked for changes in APID as well from the change server.
diff --git a/apigee_sync.go b/apigee_sync.go index e3466e3..605ef48 100644 --- a/apigee_sync.go +++ b/apigee_sync.go
@@ -1,22 +1,24 @@ package apidApigeeSync import ( + "bytes" "encoding/json" "errors" + "github.com/30x/transicator/common" "io/ioutil" "net/http" "net/url" "path" - "strconv" - "strings" "time" ) // todo: The following was largely copied from old APID - needs review -var latestMsgID int64 +var latestSequence int64 var token string -var tokenActive bool +var tokenActive, downloadSnapshot, downloadBootSnapshot, gotSequence bool +var lastSequence string +var snapshotInfo string /* * Helper function that sleeps for N seconds, if comm. with change agent @@ -55,29 +57,45 @@ */ func pollChangeAgent() error { - changesUri, err := url.Parse(config.GetString(configProxyServerBaseURI)) + if downloadSnapshot != true { + log.Error("Waiting for snapshot download to complete") + return errors.New("Snapshot download in progress...") + } + changesUri, err := url.Parse(config.GetString(configChangeServerBaseURI)) if err != nil { - log.Errorf("bad url value for config %s: %s", configProxyServerBaseURI, err) + log.Errorf("bad url value for config %s: %s", changesUri, err) return err } - changesUri.Path = path.Join(changesUri.Path, "/v1/edgex/changeagent/changes") + changesUri.Path = path.Join(changesUri.Path, "/changes") + + /* + * FIXME: This is a hack, while the correct procedure it to use the + * bootstrap scope + */ + configId := config.GetString(configScopeId) for { log.Debug("polling...") - org := config.GetString(configOrganization) - /* token not valid try again */ if tokenActive == false { - status := getTokenForOrg(org) + /* token not valid?, get a new token */ + status := getBearerToken() if status == false { return errors.New("Unable to get new token") } } + /* Find the scopes associated with the config id */ + scopes := findScopesforId(configId) /* A Blocking call for 1 Minute */ v := url.Values{} - v.Add("since", strconv.FormatInt(int64(latestMsgID), 10)) + if gotSequence == true { + v.Add("since", lastSequence) + } v.Add("block", "60") - v.Add("tag", "org:"+org) + for _, scope := range scopes { + v.Add("scope", scope) + } + v.Add("snapshot", snapshotInfo) changesUri.RawQuery = v.Encode() uri := changesUri.String() log.Info("Fetching changes: ", uri) @@ -96,12 +114,15 @@ if r.StatusCode != http.StatusOK { if r.StatusCode == http.StatusUnauthorized { tokenActive = false + log.Errorf("Token expired? Unauthorized request.") } r.Body.Close() + log.Errorf("Get Changes request failed with Resp err: %d", + r.StatusCode) return err } - var resp ChangeSet + var resp common.ChangeList err = json.NewDecoder(r.Body).Decode(&resp) r.Body.Close() if err != nil { @@ -110,73 +131,35 @@ } if len(resp.Changes) > 0 { - events.Emit(ApigeeSyncEventSelector, resp) - - lastMsgID := resp.Changes[len(resp.Changes)-1].LastMsId - if lastMsgID > 0 { - log.Infof("Updated last msg id for org %s is %s", org, lastMsgID) - - err = storeLastMsgID(org, lastMsgID) - if err != nil { - // todo: what is appropriate recovery (if anything)? - return err - } - - latestMsgID = lastMsgID - } + events.Emit(ApigeeSyncEventSelector, &resp) } else { - log.Error("Change message decoding error for org ", org) + log.Error("No Changes detected for Scopes ", scopes) } + lastSequence = resp.LastSequence + gotSequence = true } } /* - * Persist the Last Change Id in the DB + * This function will (for now) use the Access Key/Secret Key/ApidConfig Id + * to get the bearer token, and the scopes (as comma separated scope) */ -func storeLastMsgID(org string, lastID int64) error { - - db, err := data.DB() - if err != nil { - return err - } - - result, err := db.Exec("UPDATE change_id SET snapshot_change_id=? WHERE org=?;", lastID, org) - if err != nil { - log.Errorf("UPDATE change_id failed (%s: %s): %s", lastID, org, err) - return err - } - - rowsAffected, err := result.RowsAffected() - if err == nil && rowsAffected == 0 { - _, err = db.Exec("INSERT INTO change_id (snapshot_change_id, org) VALUES (?, ?);", lastID, org) - } - if err != nil { - log.Errorf("UPDATE change_id failed (%s: %s): %s", lastID, org, err) - return err - } - - log.Info("UPDATE change_id success (%s: %s)", lastID, org) - return nil -} - -func getTokenForOrg(org string) bool { +func getBearerToken() bool { uri, err := url.Parse(config.GetString(configProxyServerBaseURI)) if err != nil { log.Error(err) return false } - uri.Path = path.Join(uri.Path, "/v1/edgex/accesstoken") + uri.Path = path.Join(uri.Path, "/accesstoken") tokenActive = false form := url.Values{} - form.Add("grantType", "client_credentials") - form.Add("org", org) - req, err := http.NewRequest("POST", uri.String(), strings.NewReader(form.Encode())) + form.Set("grantType", "client_credentials") + form.Add("client_id", config.GetString(configConsumerKey)) + form.Add("client_secret", config.GetString(configConsumerSecret)) + req, err := http.NewRequest("POST", uri.String(), bytes.NewBufferString(form.Encode())) req.Header.Set("Content-Type", "application/x-www-form-urlencoded; param=value") - - consumerKey := config.GetString(configConsumerKey) - req.SetBasicAuth(consumerKey, config.GetString(configConsumerSecret)) client := &http.Client{} resp, err := client.Do(req) if err != nil { @@ -202,7 +185,7 @@ } token = oauthResp.AccessToken tokenActive = true - log.Info("Got a new token for Consumer: ", consumerKey) + log.Info("Got a new token..") return true } @@ -223,26 +206,38 @@ func Redirect(req *http.Request, via []*http.Request) error { req.Header.Add("Authorization", "Bearer "+token) - req.Header.Add("org", config.GetString(configOrganization)) + req.Header.Add("org", config.GetString(configScopeId)) return nil } -func downloadSnapshot() error { +func DownloadSnapshot() error { - org := config.GetString(configOrganization) - status := getTokenForOrg(org) +RETRY: + var scopes []string + + /* Get the bearer token */ + status := getBearerToken() if status == false { return errors.New("Unable to get new token") } - snapshotUri, err := url.Parse(config.GetString(configProxyServerBaseURI)) + snapshotUri, err := url.Parse(config.GetString(configSnapServerBaseURI)) if err != nil { - log.Errorf("bad url value for config %s: %s", configProxyServerBaseURI, err) - return err + log.Fatalf("bad url value for config %s: %s", snapshotUri, err) } - snapshotUri.Path = path.Join(snapshotUri.Path, "/v1/edgex/snapshots?org") + + if downloadBootSnapshot == false { + scopes = append(scopes, (config.GetString(configScopeId))) + } else { + scopes = findScopesforId(config.GetString(configScopeId)) + } + + /* Frame and send the snapshot request */ + snapshotUri.Path = path.Join(snapshotUri.Path, "/snapshots") v := url.Values{} - v.Add("tag", org) + for _, scope := range scopes { + v.Add("scopes", scope) + } snapshotUri.RawQuery = v.Encode() uri := snapshotUri.String() log.Info("Snapshot Download : ", uri) @@ -252,29 +247,60 @@ } req, err := http.NewRequest("GET", uri, nil) req.Header.Add("Authorization", "Bearer "+token) - resp, err := client.Do(req) + r, err := client.Do(req) if err != nil { - log.Error("API Proxy comm error: [%s] ", err) + log.Fatalf("Snapshotserver comm error: [%s] ", err) + } + defer r.Body.Close() + + /* Decode the Snapshot server response */ + var resp common.Snapshot + err = json.NewDecoder(r.Body).Decode(&resp) + if err != nil { + log.Fatalf("JSON Response Data not parsable: [%s] ", err) return err } - defer resp.Body.Close() - if resp.StatusCode == 200 { - rawjson, err := ioutil.ReadAll(resp.Body) - if err != nil { - log.Error("Snapshot response read error: [%s] ", err) - return err - } - err = data.InsertSnapshotDB(rawjson) - if err != nil { - log.Error("Insert Snapshot error: [%s] ", err) - return err + /* + * The idea here is that you download snapshot for the scopes + * associated with the apidconfig Id, and then download the + * data based on the scopes retrieved in the first round + */ + if r.StatusCode == 200 { + log.Info("Emit Snapshot response to plugins") + events.Emit(ApigeeSyncEventSelector, &resp) + snapshotInfo = resp.SnapshotInfo + if downloadBootSnapshot == false { + downloadBootSnapshot = true + goto RETRY + } else if downloadBootSnapshot == true { + downloadSnapshot = true } - - log.Info("Got a new DB from Snapshot server") - return err + } else { + log.Fatalf("Snapshot server Connect failed. Resp code %d", r.StatusCode) } - log.Info("Snapshot server Connect failed. Resp code %d", resp.StatusCode) + return err +} +func findScopesforId(configId string) (scopes []string) { + + var scope string + db, err := data.DB() + if err != nil { + log.Errorf("DB open Error: %s", err) + return nil + } + + rows, err := db.Query("select scope from APID_CONFIG_SCOPE where apid_config_id = $1", configId) + if err != nil { + log.Errorf("Failed to query APID_CONFIG_SCOPE. Err: %s", err) + return nil + } + defer rows.Close() + for rows.Next() { + rows.Scan(&scope) + scopes = append(scopes, scope) + } + return scopes }
diff --git a/init.go b/init.go index e24acd4..81bc050 100644 --- a/init.go +++ b/init.go
@@ -1,16 +1,19 @@ package apidApigeeSync import ( + "database/sql" "fmt" "github.com/30x/apid" ) const ( - configOrganization = "apigeesync_organization" // todo: how are we supporting multiple orgs? - configPollInterval = "apigeesync_poll_interval" - configProxyServerBaseURI = "apigeesync_proxy_server_base" - configConsumerKey = "apigeesync_consumer_key" - configConsumerSecret = "apigeesync_consumer_secret" + configPollInterval = "apigeesync_poll_interval" + configProxyServerBaseURI = "apigeesync_proxy_server_base" + configSnapServerBaseURI = "apigeesync_snapshot_server_base" + configChangeServerBaseURI = "apigeesync_change_server_base" + configConsumerKey = "apigeesync_consumer_key" + configConsumerSecret = "apigeesync_consumer_secret" + configScopeId = "apigeesync_bootstrap_id" ApigeeSyncEventSelector = "ApigeeSync" ) @@ -36,23 +39,68 @@ config.SetDefault(configPollInterval, 120) + db, err := data.DB() + if err != nil { + log.Panic("Unable to access DB", err) + } + // check for required values - for _, key := range []string{configProxyServerBaseURI, configOrganization, configConsumerKey, configConsumerSecret} { + for _, key := range []string{configProxyServerBaseURI, configConsumerKey, configConsumerSecret, configSnapServerBaseURI, configChangeServerBaseURI} { if !config.IsSet(key) { return fmt.Errorf("Missing required config value: %s", key) } } - err := downloadSnapshot() - if err != nil { - log.Error("Unable to download snapshot") - return nil + var count int + row := db.QueryRow("SELECT count(*) FROM sqlite_master WHERE type='table' AND name='apid_config' COLLATE NOCASE;") + if err := row.Scan(&count); err != nil { + log.Panic("Unable to setup database", err) + } + if count == 0 { + createTables(db) } + /* call to Download Snapshot info */ + go DownloadSnapshot() + + /* Begin Looking for changes periodically */ log.Debug("starting update goroutine") go updatePeriodicChanges() + events.Listen(ApigeeSyncEventSelector, &handler{}) + log.Debug("end init") return nil } + +func createTables(db *sql.DB) { + _, err := db.Exec(` +CREATE TABLE apid_config ( + id text, + name text, + description text, + umbrella_org_app_name text, + created int64, + created_by text, + updated int64, + updated_by text, + _apid_scope text, + PRIMARY KEY (id) +); +CREATE TABLE apid_config_scope ( + id text, + apid_config_id text, + scope text, + created int64, + created_by text, + updated int64, + updated_by text, + _apid_scope text, + PRIMARY KEY (id) +); +`) + if err != nil { + log.Panic("Unable to initialize DB", err) + } +}
diff --git a/listener.go b/listener.go new file mode 100644 index 0000000..3730e14 --- /dev/null +++ b/listener.go
@@ -0,0 +1,163 @@ +package apidApigeeSync + +import ( + "database/sql" + "github.com/30x/apid" + "github.com/30x/transicator/common" +) + +type handler struct { +} + +func (h *handler) String() string { + return "ApigeeSync" +} + +// todo: The following was basically just copied from old APID - needs review. + +func (h *handler) Handle(e apid.Event) { + + snapData, ok := e.(*common.Snapshot) + if ok { + processSnapshot(snapData) + } else { + changeSet, ok := e.(*common.ChangeList) + if ok { + processChange(changeSet) + } else { + log.Errorf("Received Invalid event. This shouldn't happen!") + } + } + return +} + +func processSnapshot(snapshot *common.Snapshot) { + + log.Debugf("Process Snapshot data") + + db, err := data.DB() + if err != nil { + panic("Unable to access Sqlite DB") + } + + for _, payload := range snapshot.Tables { + + switch payload.Name { + case "apid_config": + for _, row := range payload.Rows { + insertApidConfig(row, db) + } + case "apid_config_scope": + for _, row := range payload.Rows { + insertApidConfigScope(row, db) + } + } + } +} + +func processChange(changes *common.ChangeList) { + + log.Debugf("apigeeSyncEvent: %d changes", len(changes.Changes)) + + db, err := data.DB() + if err != nil { + panic("Unable to access Sqlite DB") + } + + for _, payload := range changes.Changes { + + switch payload.Table { + case "public.apid_config": + case "edgex.apid_config": + switch payload.Operation { + case 1: + insertApidConfig(payload.NewRow, db) + } + + case "public.apid_config_scope": + case "edgex.apid_config_scope": + switch payload.Operation { + case 1: + insertApidConfigScope(payload.NewRow, db) + } + } + } +} + +/* + * INSERT INTO APP_CREDENTIAL op + */ +func insertApidConfig(ele common.Row, db *sql.DB) { + + var scope, id, name, orgAppName, createdBy, updatedBy, Description string + var updated, created int64 + + txn, _ := db.Begin() + err := ele.Get("id", &id) + err = ele.Get("_apid_scope", &scope) + err = ele.Get("name", &name) + err = ele.Get("umbrella_org_app_name", &orgAppName) + err = ele.Get("created", &created) + err = ele.Get("created_by", &createdBy) + err = ele.Get("updated", &updated) + err = ele.Get("updated_by", &updatedBy) + err = ele.Get("description", &Description) + + _, err = txn.Exec("INSERT INTO apid_config (id, _apid_scope, name, umbrella_org_app_name, created, created_by, updated, updated_by)VALUES(?,?,?,?,?,?,?,?);", + id, + scope, + name, + orgAppName, + created, + createdBy, + updated, + updatedBy, + Description) + + if err != nil { + log.Error("INSERT Failed: ", id, ", ", scope, ")", err) + txn.Rollback() + } else { + log.Info("INSERT Success: (", id, ", ", scope, ")") + txn.Commit() + } + +} + +/* + * INSERT INTO APP_CREDENTIAL op + */ +func insertApidConfigScope(ele common.Row, db *sql.DB) { + + var id, scopeId, apiConfigId, scope, createdBy, updatedBy string + var created, updated int64 + + txn, _ := db.Begin() + err := ele.Get("id", &id) + err = ele.Get("_apid_scope", &scopeId) + err = ele.Get("apid_config_id", &apiConfigId) + err = ele.Get("scope", &scope) + err = ele.Get("created", &created) + err = ele.Get("created_by", &createdBy) + err = ele.Get("updated", &updated) + err = ele.Get("updated_by", &updatedBy) + + _, err = txn.Exec("INSERT INTO apid_config_scope (id, _apid_scope, apid_config_id, scope, created, created_by, updated, updated_by)VALUES(?,?,?,?,?,?,?,?);", + id, + scopeId, + apiConfigId, + scope, + created, + createdBy, + updated, + updatedBy) + + if err != nil { + log.Error("INSERT CRED Failed: ", id, ", ", scope, ")", err) + txn.Rollback() + } else { + log.Info("INSERT CRED Success: (", id, ", ", scope, ")") + txn.Commit() + } + +}