Wednesday, October 16, 2019

Twitter Analytics with Google Natural Language


Summary

I'll be taking Twitter tweets and processing them through Google's Natural Language APIs in this post.  The NL APIs provide the ability to parse text into 'entities' and/or determining 'sentiment' of the entity and surrounding text.  I'll be using a combination of the two to analyze some tweets.

Entity + Sentiment Analysis

Lines 4-23:  REST call to Google's NL entity-sentiment endpoint.  The response from that endpoint is an array of entities in rank order of 'salience' (relevance).  Entities have types, such as organization, person, etc.  Net, the tweet gets parsed into entities and I'm pulling the #1 entity from the tweet as ranked by salience.

const ENTITY_SENTIMENT_URL = 'https://language.googleapis.com/v1beta2/documents:analyzeEntitySentiment?key=' + GOOGLE_KEY;

    try {
        const response = await fetch(ENTITY_SENTIMENT_URL, {
            method : 'POST',
            body : JSON.stringify(body),
            headers: {'Content-Type' : 'application/json; charset=utf-8'},
        });

        if (response.ok) {
            const json = await response.json();
            const topSalience = json.entities[0];
            const results = {
                'name' : topSalience.name,
                'type' : topSalience.type,
                'salience' : topSalience.salience,
                'entitySentiment' : topSalience.sentiment
            }
            return results;
        }
        else {
            let msg = (`response status: ${response.status}`);
            throw new Error(msg);
        }
    }
    catch (err) {
        ts = new Date();
        let msg = (`${ts.toISOString()} entitySentiment() - ${err}`);
        console.error(msg)
        throw err;
    }

Sentiment Analysis

Lines 1-22 implement a sentiment analysis of the entire tweet text.
    try {
        const response = await fetch(SENTIMENT_URL, {
            method : 'POST',
            body : JSON.stringify(body),
            headers: {'Content-Type' : 'application/json; charset=utf-8'},
        });

        if (response.ok) {
            const json = await response.json();
            return json.documentSentiment;
        }
        else {
            let msg = (`response status: ${response.status}`);
            throw new Error(msg);
        }
    }
    catch (err) {
        ts = new Date();
        let msg = (`${ts.toISOString()} sentiment() - ${err}`);
        console.error(msg)
        throw err;
    }

Blending

Lines 1-16:  This function calls upon both of the Google NL API functions above and provides a blended analysis of the tweet text.  Below, I took the product of the averages of the magnitude (amount of emotion) and score (positive vs negative emotion) between the entity-sentiment and overall sentiment to arrive at an aggregate figure.  There are certainly other ways to combine these factors.

async function analyze(tweet) {
    const esnt = await entitySentiment(tweet);
    const snt = await sentiment(tweet);

    let results = {};
    results.tweet = tweet;
    results.name = esnt.name;
    results.type = esnt.type;
    results.salience = esnt.salience;
    results.entitySentiment = esnt.entitySentiment;
    results.documentSentiment = snt;
    let mag = (results.entitySentiment.magnitude + results.documentSentiment.magnitude) / 2;
    let score = (results.entitySentiment.score + results.documentSentiment.score) / 2;
    results.aggregate = mag * score;
    return results;
}

Execution

Lines 1-4:  Simple function for reading a JSON-formatted file of tweets.
Lines 6-17:  Reads a file containing an array of tweets and process each through the Google NL functions mentioned above.
async function readTweetFile(file) {
    let tweets = await fsp.readFile(file);
    return JSON.parse(tweets);
}

readTweetFile(INFILE)
.then(tweets => {
    for (let i=0; i < tweets.length; i++) {
        analyze(tweets[i].text)
        .then(json => {
            console.log(JSON.stringify(json, null, 4));
        });
    }
})
.catch(err => {
    console.error(err);
});

Results

Below is an example of a solid negative sentiment from Donald Trump regarding Adam Schiff.  Schiff is accurately identified as a 'Person' entity.  Note the negative scores + high emotion (magnitude) in both the entity sentiment and overall sentiment analysis.

{
    "tweet": "Shifty Adam Schiff wants to rest his entire case on a Whistleblower who he now
     says can’t testify, & the reason he can’t testify is that he is afraid to do so 
     because his account of the Presidential telephone call is a fraud & 
     totally different from the actual transcribed call...",
    "name": "Adam Schiff",
    "type": "PERSON",
    "salience": 0.6048015,
    "entitySentiment": {
        "magnitude": 3.2,
        "score": -0.5
    },
    "documentSentiment": {
        "magnitude": 0.9,
        "score": -0.9
    },
    "aggregate": -1.435
}

Below is another accurately detected 'person' entity with a positive statement from the President.
{
    "tweet": "Kevin McAleenan has done an outstanding job as Acting Secretary of
Homeland Security. We have worked well together with Border Crossings being way down.
Kevin now, after many years in Government, wants to spend more time with his family and
go to the private sector....",
    "name": "Kevin McAleenan",
    "type": "PERSON",
    "salience": 0.6058554,
    "entitySentiment": {
        "magnitude": 0.4,
        "score": 0
    },
    "documentSentiment": {
        "magnitude": 1.2,
        "score": 0.3
    },
    "aggregate": 0.12
}

Source

https://github.com/joeywhelan/twitterAnalytics

Copyright ©1993-2024 Joey E Whelan, All rights reserved.

Sunday, October 13, 2019

Twitter Premium Search API - Node.js


Summary

In this post I'll demonstrate how to use the Twitter Premium Search API.   This is a pure REST API with two different search modes:  past 30 days or full archive search since Twitter existed (2006).

The API has a very limited 'free' mode for Developers to try out.  Limits are imposed on usage:  number of API requests, tweets pulled per month and rate of API calls.  To do anything of significance with this API, you're faced with paying for Twitter's API subscription.  That gets pretty pricey quickly with the cheapest tier currently at $99/month.  This post is based on usage of the 'free'/sandbox tier.

Main Loop

Line 1 fetches a bearer token for accessing the Twitter APIs.  I covered this topic in a previous post.

Lines 4-19 implement a while loop that fetches batches of tweets for a given search query.  For the Twitter free/sandbox environment, you can pull up to 100 tweets per API call.  Each tweet in the batch is evaluated to determine if it was 140 or 280 character tweet.  The tweet text is formatted and then that and the created_date are added to a JSON array.  That array is ultimately written to file.

Line 20 is a self-imposed delay on calls to the Twitter API.  If you bust their rate limits, you'll get a HTTP 429 error.

        const token = await getTwitterToken(AUTH_URL);
        let next = null;
        
        do {
            const batch = await getTweetBatch(token, url, query, fromDate, maxResults, next);
            for (let i=0; i < batch.results.length; i++) {  //loop through the page/batch of results
                let tweet = {};
                if (batch.results[i].truncated) {  //determine if this is a 140 or 280 character tweet
                    tweet.text = batch.results[i].extended_tweet.full_text.trim();
                }
                else {
                    tweet.text = batch.results[i].text.trim();
                }

                tweet.text = tweet.text.replace(/\r?\n|\r|@|#/g, ' ');  //remove newlines, @ and # from tweet text
                tweet.created_at = batch.results[i].created_at;
                tweets.push(tweet);
            }
            next = batch.next;
            await rateLimiter(3);  //rate limit twitter api calls to 1 per 3 seconds/20 per minute
        }
        while (next);

Tweet Batch Fetch

Lines 1-26 set up a node fetch to the Twitter REST API end point.  If this was a call with a 'next' parameter (meaning multiple pages of tweets on a single search), I add that parameter to the fetch.

    const body = {
        'query' : query,
        'fromDate' : fromDate,
        'maxResults' : maxResults
    };
    if (next) {
        body.next = next;
    }

    try {
        const response = await fetch(url, {
            method: 'POST',
            headers: {
            'Authorization' : 'Bearer ' + token
            },
            body: JSON.stringify(body)
        });
        if (response.ok) {
            const json = await response.json();
            return json;
        }
        else {
            let msg = (`authorization request response status: ${response.status}`);
            throw new Error(msg);    
        }
    }

Usage

let query = 'from:realDonaldTrump -RT';  //get tweets originated from Donald Trump, filter out his retweets
let url = SEARCH_URL + THIRTY_DAY_LABEL;  //30day search
let fromDate = '201910010000'; //search for tweets within the current month (currently, Oct 2019)
search(url, query, fromDate, 100)  //100 is the max results per request for the sandbox environment 
.then(total => {
    console.log('total tweets: ' + total);
})
.catch(err => {
    console.error(err);
});

Output

Snippet of the resulting JSON array from the function call above.
[
    {
        "text": "We have become a far greater Economic Power than ever before, and we are using that power for WORLD PEACE!",
        "created_at": "Sun Oct 13 14:32:37 +0000 2019"
    },
    {
        "text": "Where’s Hunter? He has totally disappeared! Now looks like he has raided and scammed even more countries! Media is AWOL.",
        "created_at": "Sun Oct 13 14:15:55 +0000 2019"
    },

Source

https://github.com/joeywhelan/twitterSearch

Copyright ©1993-2024 Joey E Whelan, All rights reserved.

Sunday, September 22, 2019

Fetching an Application-only Twitter API Token via Node


Summary


This short post demonstrates the steps necessary to fetch an app-only bearer token via Twitter's OAuth2 interface using Node.js.  That token would then be subsequently used to access Twitter's APIs.  This post follows the steps explained on the Twitter developer site here.

Set-up

  1. Create a developer account as described here.  
  2. Create an 'application' as described here.
  3. At this point, you have a 'Consumer Key' and 'Consumer Secret'.  Those two strings will be used in the code shown below.

Code

 

Create the Consumer Token

Per the Twitter documentation, the Consumer Key and Secret need to be URL encoded, concatentated, and then base64-encoded.
const CONSUMER_KEY = process.env.CONSUMER_KEY;
const CONSUMER_SECRET = process.env.CONSUMER_SECRET;

function urlEncode (str) {
    return encodeURIComponent(str)
      .replace(/!/g, '%21')
      .replace(/'/g, '%27')
      .replace(/\(/g, '%28')
      .replace(/\)/g, '%29')
      .replace(/\*/g, '%2A')
}

const consumerToken = btoa(urlEncode(CONSUMER_KEY) + ':' + urlEncode(CONSUMER_SECRET));

Fetch the Bearer Token

Code below uses node-fetch to execute a HTTP POST to the Twitter OAuth2 interface.  If the fetch is successful, the bearer token is inside a JSON object returned by that interface.

    return fetch(url, {
        method: 'POST',
        headers: {
            'Authorization' : 'Basic ' + consumerToken,
            'Content-Type' : 'application/x-www-form-urlencoded;charset=UTF-8'
        }, 
        body : 'grant_type=client_credentials'
    })
    .then(response => {
        if (response.ok) {
            return response.json();
        }
        else {
            throw new Error('Response Status: ' + response.status);
        }
    })
    .then(json => {
        if (json.token_type == 'bearer') {
            return json.access_token;
        }
        else {
            throw new Error('Invalid token type: ' + json.token_type);
        }
    });  

Source


Full source here: https://github.com/joeywhelan/authTest

Copyright ©1993-2024 Joey E Whelan, All rights reserved.

Sunday, July 21, 2019

API Development with GCP


Summary

This post covers my use of Google Cloud Platform (GCP) to build a scalable piece of middleware, theoretically, infinitely scalable.  The notional use case here is a store locator microservice for a company that has thousands of locations.  The service will provide the closest store location for a given ZIP code or GPS coordinates.  The location will be returned as a URL representing the directions on Google Maps.

Overall Architecture

For this particular application, I utilized GCP-based services for the entirety of the core app.  I additionally leveraged the services of Auth0 for authentication services.  The API is realized in a REST structure - two GET-based services.

GCP Architecture

I used Google App Engine Flex (GAE) for the application core.  I provided front-end API support with Cloud Endpoints.  Caching of store and  ZIP coordinates utilizes Google's cloud implementation of Redis - Cloud Memorystore.  A combination of Cloud Functions and Cloud Storage provides the ability for real-time updating of the cache with simple configuration file modifications.

API Proxy + Authentication Layer

All requests into the webserver written in node.js on GAE Flexible are proxied by Cloud Endpoints.  Endpoints provides redirection to HTTPS and authentication proxying, among other things.  I chose utilize the OAuth services provided by Auth0.  They in turn proxy the JWT-based (RSA256) authentication that provides security for the API.



Configuration of the Cloud Endpoints is accomplished via an OpenAPI ver 2.0 YAML file.

swagger: "2.0"
info:
  title: "Store Locator"
  description: "Generate a Google maps directions URL to the nearest 
  store based on user's current ZIP code or coordinates"
  version: "1.0.0"
host: "youraccount.appspot.com"
consumes:
- "text/plain"
produces:
- "text/uri-list"
schemes:
  - "https"
securityDefinitions:
  auth0_jwk:
    authorizationUrl: ""
    flow: "implicit"
    type: "oauth2"
    x-google-issuer: "https://yourcaccount.auth0.com/"
    x-google-jwks_uri: "https://youraccount.auth0.com/.well-known/jwks.json"
    x-google-audiences: "https://locatorUser1"
security:
  - auth0_jwk: []
paths:
  /locator/zip:
    get:
      summary: "Find nearest store by ZIP code"
      operationId: "ZIP"
      description: "ZIP code"
      parameters:
        -
          name: zip
          in: query
          required: true
          type: string
      responses:
        200:
          description: "Google Maps URL with directions from input ZIP to nearest store"
          schema: 
            type: string
        404:
          description: "Error Message"
          schema:
            type: string
  /locator/coordinates:
    get:
      summary: "Find nearest store by coordinates (latitude, longitude)"
      operationId: "Coordinates"
      description: "Latitude, Longitude"
      parameters:
        -
          name: coordinates
          in: query
          required: true
          type: string
      responses:
        200:
          description: "Google Maps URL with directions from input coordinates to nearest store"
          schema: 
            type: string
        404:
          description: "Error Message"
          schema:
            type: string


Application Layer

The heart of this is a node.js application that realizes two HTTP GET paths in Express.  The paths allow for a search of the closest store based on the user's current ZIP code or GPS coordinates (latitude + longitude).  Cloud Memorystore keeps an updated set of coordinates for the ZIP codes and store locations.

Code Snippet - Get Closest Store by Coordinates

I use a filter based on the haversine formula to narrow down the list of closest store candidates.  That formula finds 'as the crow flies' distances.  Once those candidates have been found, I then send them into a Google Maps API service (Distance Matrix) for actual driving distances.  Ultimately, the API call returns a URL that corresponds to the actual driving directions between the origin and store coordinates.

/**
 * Performs the Haversine formula to generate the great circle distance between two coordinates.
 * @param {object} coord1 - latitude & longitude
 * @param {object} coord2 - latitude & longitude
 * @return {int} - great circle distance between the two coordinates
 */
function haversine(coord1, coord2) {
 let lat1 = coord1.lat;
 let lon1 = coord1.long;
 let lat2 = coord2.lat;
 let lon2 = coord2.long;
 const R = 3961;  //miles
    const degRad = Math.PI/180;
    const dLat = (lat2-lat1)*degRad;
    const dLon = (lon2-lon1)*degRad;
 
 lat1 *= degRad;
    lat2 *= degRad;
    const a = Math.sin(dLat/2) * Math.sin(dLat/2) + 
        Math.sin(dLon/2) * Math.sin(dLon/2) * Math.cos(lat1) * Math.cos(lat2);
    const c = 2 * Math.atan2(Math.sqrt(a), Math.sqrt(1-a));
    return R * c;
}

/**
 * Fetches a configurable number of stores that are closest to a given coordinate.
 * @param {object} origin - latitude & longitude
 * @param {int} numVals - number of closest stores to return
 * @return {array} - array of the numVal closest stores
 */
function getStoresByCoord(origin, numVals) {
 console.log(`getStoresByCoord(JSON.stringify(origin), numVals)`);
 let distances = [];
 
 if (numVals > storeList.length) numVals = storeList.length;

 //performs a haversine dist calc between the origin and each of the stores
 for (let i=0; i < storeList.length; i++) {
  const dist = haversine(origin, {'lat' : storeList[i].lat, 'long' : storeList[i].long});
  const val = {'index': i, 'distance': dist};
  distances.push(val);
 }

 let stores = [];
 distances.sort(compareDist);
 for (let i = 0; i < numVals; i++) {
  stores.push(storeList[distances[i].index]);
 }
 
 return stores;
}

/**
 * Fetches the closest store location based on the user's lat/long
 * Performs an initial filtering based ZIP code only, then refines a configurable number of closest
 * stores using actual driving distance from the Google Maps Distance Matrix
 * API.
 * @param {string} origin - lat/long of origin location
 * @param {array of strings} stores - lat/long(s) of store locations
 * @return {string} - URL of Google map with nearest store
 */
function getClosestStore(origin, stores) {
 console.log(`getClosestStore(JSON.stringify(origin))`);
 
 let dests = [];
 for (const store of stores) {
  dests.push({lat: store.lat, lng: store.long});
 }
 return mapsClient.distanceMatrix({
  'origins': [{lat: origin.lat, lng: origin.long}],
  'destinations': dests
 })
 .asPromise()
 .then((response) => {
  if (response.status == 200 && response.json.status == 'OK') {
   let minIndex;
   let minDist;
   const elements = response.json.rows[0].elements
   for (let i=0; i < elements.length; i++) {
    if (minDist == null || elements[i].distance.value < minDist) {
     minIndex = i;
     minDist = elements[i].distance.value;
    }
   }
   return stores[minIndex];
  }
  else {
   throw new Error('invalid return status on Google distanceMatrix API call');
  }
 })
 .catch(err => {
  console.error(`getClosestStore(): ${JSON.stringify(err)}`);
  throw err;
 });
}

/**
 * Creates a google maps url with the directions from an origin to a store location.
 * @param {object} origin - latitude & longitude
 * @param {object} store - object containing store address info, to include latitude & longitude
 * @return {string} - Google Maps URL showing directions from origin to store location
 */
function getDirections(origin, store) {
 return MAPSURL + `&origin=origin.lat, origin.long` + `&destination=store.lat, store.long`;
}

/**
 * Fetches the closest store location based on the user's lat/long.
 * @param {string} coordinates - lat/long of user's current location
 * @return {string} - URL of Google map with nearest store
 */
app.get('/locator/coordinates', (request, response) => {
 const vals = request.query.coordinates.split(',');
 const origin = {'lat' : vals[0], 'long' : vals[1]};
 const stores = getStoresByCoord(origin, 3);
 getClosestStore(origin, stores)
 .then(store => {
  const url = getDirections(origin, store);
  response.status(200).send(url);
 })
 .catch(err => {
  response.status(404).send(err.message);
 });
});

Caching Layer

Google Cloud Memorystore is a cloud implementation of Redis.  I use this service to provide a real-time updatable set of ZIP code and store location coordinates.

Configuration + Persistence Layer

This app is able to get to real-time updates on stores and ZIP codes simply by reloading files into Google Cloud Storage.  This allows move/add/deletes of stores without any modification of code or restarts of the application.  I created a Cloud Function that is triggered on modification of either the stores or ZIP code files.  After being triggered, the Cloud Function updates the Memorystore cache with the latest information.

Code Snippet - Cloud Function (gcsMonitor.js)


/**
* Public function for reading the Store location file from Google Cloud Storage
* and loading into Google CloudMemory(redis).  
* Will propagate exceptions.
*/
function loadStoreCache() {
 console.log(`loadStoreCache() executing`);
 const bucket = storage.bucket(gcpBucket);
 const stream = bucket.file(gcpStoreFile).createReadStream();
 const client = redis.createClient(REDISPORT, REDISHOST);
 client.on("error", function (err) {
  console.log("loadStoreCache() Redis error:" + err);
 });

 csv()
 .fromStream(stream)
 .subscribe((json) => {
  let hashKey;
  for (let [key, value] of Object.entries(json)) {
   if (key === 'storeNum') {
    hashKey = 'store:' + value;
   }
   else {
    console.log(`loadStoreCache() inserting hashKey`);
    client.hset(hashKey, key, value, (err, reply) => {
     if (err) {
      console.error(`loadStoreCache() Error: ${err}`);
     }
     
    });
   }
  }
 })
 .on('done', (err) => {
  client.quit();
  console.log(`loadStoreCache() complete`);
 });
};

Test Client

As discussed, the API uses Auth0 (proxied by Endpoints) for authentication.  The code below fetches a JWT token from Auth0 and then performs an API call with that token in the header.
function fetchToken() {

    const body = {
        'client_id': clientId,
        'client_secret': clientSecret,
        'audience': audience,
        'grant_type': 'client_credentials'
    };

    return fetch(tokenUrl, {
        method: 'POST',
        headers: {
            'Content-Type': 'application/json'
        },
        body: JSON.stringify(body)
    })
    .then(response => {
        if (response.ok) {
            return response.json();
        }
        else {
            console.error('fetchToken Error: ' + response.status);
        }
    }) 
    .then(json => {
        return json.access_token;
    })   
}

const apiUrl = 'https://yourapp.appspot.com/locator/coordinates/?coordinates=37.1464,-94.4630'
function authTest() {
    return fetchToken()
    .then(token => {
        fetch(apiUrl, {
            method: 'GET',
            headers: {
                'Authorization': 'Bearer ' + token
            }
        })
        .then(response => {
            if (response.ok) {
                return response.text();
            }
            else {
                console.error(response.status);
            }
        })
        .then(text => {
            console.log('Response: ' + text);
        })
    });
}

authTest();


Results

$ node testClient.js
Response: https://www.google.com/maps/dir/?api=1&origin=37.1464, -94.4630&destination=37.0885389, -94.5144897

Source


Copyright ©1993-2024 Joey E Whelan, All rights reserved.

Saturday, March 2, 2019

Event Sourcing with Redis Streams


Summary

Streams were an addition to the Redis 5.0 release.  Redis streams are roughly analogous to a log file: an append-only data structure.  Redis streams also have some producer/consumer capabilities that are roughly analogous to Kafka streams, with some important differences.  A full explanation of Redis streams here.

Event sourcing is one of those old topics that has been brought back to life again in a new context.  The new context is state persistence and messaging for microservice architectures.  A full explanation of event sourcing here.

This post is about my adventures at implementing event sourcing via Redis streams.  The example is simple and contrived - an account microservice that allows deposits and withdrawals.  I made a rough attempt at implementing this in a 'domain-driven design' model, but I didn't fully adhere to that model as I didn't care for the level of abstraction necessary.

Even though the scenario is simple, my overall impression of event sourcing is - it's hard.  It's hard to think about code in an event-driven manner and it's hard to implement it correctly in a distributed architecture.

Overall Architecture

Diagram below of the high-level architecture:  REST-based microservice that leverages Redis for the event store and MongoDB for event data aggregation.

Service Architecture

Microserver application arch below.  I implemented this with a Node.js HTTP server for the REST routes and an Account Service class that has an event store client and an array of account aggregates.

Projection Architecture

Architecture for the aggregating data from the event store below.  Again, Node.js implementation with a Redis client that realizes an event store functionality and a MongoDB client for aggregating data from events.

Event Store Architecture

A rough outline of what the event store implementation looks like.  I use Redis objects, in particular, the stream object to realize event sourcing functionality such as fetching an event, publishing an event, subscribing for events, etc.

Code Snippets

Creating an Account - accountService.js

The code below leverages a Redis Set object to ensure unique account ID's.  A JSON object is then created with the corresponding 'create' event and then 'published' to Redis.

 create(id) {
  return this._client.addId(id, 'accountId')
  .then(isUnique => {
   logger.debug(`AccountService.create - id:${id} - isUnique:${isUnique}`);
   if (isUnique) {
    const newEvent = {'id' : id, 'version': 0, 'type': 'create'};
    return this._client.publish('accountStream', newEvent);
   }
   else {
    return new Promise((resolve, reject) => {
     resolve(null);
    });
   }
  })
  .then (results => {
   if (results && results.length === 2) {  //results is an array.  first item is the new version number of the aggregate, 
             //second is the timestamp of the create event that was published
    logger.debug(`AccountService.create - id:${id} - results:${results}`);
    const version = results[0];
    const timestamp = results[1];
    const account = new Account(id, version, timestamp);
    this._accounts[id] = account;  //add the new account to the cache
    return {'id' : id};
   }
   else {
    throw new Error('Attempting to create an account id that already exists');
   }
  })
  .catch(err => {
   logger.error(`AccountService.create - id:${id} - ${err}`);
   throw err;
  });
 }

Making a deposit to an Account - accountService.js

The code below attempts to load the account from cache or replay of events if not in cache.  Business logic for a deposit is implemented in the account aggregate (account.js).  If the aggregate allows the deposit, then an event is published.  If the publishing of the event fails (concurrency conflict), the deposit is rolled back and an error is thrown.

 deposit(id, amount) {
  let account;
  
  return this._loadAccount(id) //attempt to load the account from cache and/or rehydrate from events
  .then(result => {
   account = result;
   account.deposit(amount);
   const newEvent = {'id' : id, 'version' : account.version, 'type': 'deposit', 'amount': amount};
   return this._client.publish('accountStream', newEvent);
  })
  .then(results => {
   logger.debug(`AccountService.deposit - id:${id}, amount:${amount} - results:${results}`);
   if (results) {
    account.version = results[0];
    account.timestamp = results[1];
    this._accounts[id] = account; //update the account cache
    return {'id': id, 'amount': amount};
   }
   else {
    account.withdraw(amount); //rolling back aggregate due to unsuccessful publishing of deposit event
    return null;
   }
  })
  .catch(err => {
   logger.error(`PlayerService.deposit - id:${id}, amount:${amount} - ${err}`);
   throw err;
  });
 }

Publishing an event - eventStoreClient.js

The code below implements concurrency control to the event store with Redis' 'watch' method.  Only one process will be permitted to publish an event with a given 'version' number.

 publish(streamName, event) { 
  logger.debug(`EventStoreClient.publish`);
  this._client.watch(event.id); //watch the id (account)
  return this._getAsync(event.id)  //fetch the current version from a Redis key with that ID 
  .then(result => { 
   if (!result || parseInt(result) === parseInt(event.version)) {  //key doesn't exist or versions match
    event.version += 1;  //increment version number prior to publishing the event
    logger.debug(`EventStoreClient.publish - streamName:${streamName}, event:${JSON.stringify(event)}\
     - result:${result}`);
    return new Promise((resolve, reject) => {
     this._client.multi()  //atomic transaction that increments the version and adds event to stream
     .incr(event.id)
     .xadd(streamName, '*', 'event', JSON.stringify(event))
     .exec((err, replies) => {
      if (err) {
       reject(err);
      }
      else {
       resolve(replies);
      }
     });
    });
   }
   else {  //covers the scenario where a concurrent access causes a mismatch with event version numbers
     //return null and then it's up to the client to make another publish attempt
    return new Promise((resolve, reject) => {
     resolve(null);
    });
   }
  })
  .catch(err => {
   logger.error(`EventStoreClient.publish - streamName:${streamName}, event:${event} - ${err}`);
   throw err;
  });
 }

Subscribing for events - eventStoreClient.js

The Redis streams implementation doesn't provide standard pub/sub functionality.  Subscriber behavior can be emulated though using Node.js event emitters coupled with Redis consumer groups.  A given consumer group is read below periodically via setInterval.  If new events are present, they're emitted via the emitter.  The 'subscriber' would implement the corresponding Node event handler for the emitter returned by this function.  This is precisely how the 'projector' is implemented for aggregating event data to a MongoDB database.

 subscribe(streamName, consumerName) { 
  let emitter;
  let groupName = streamName + 'Group';
  logger.debug(`EventStoreClient.subscribe - streamName:${streamName}, groupName:${groupName}, consumerName:${consumerName}`);
  
  if (this._emitters[streamName] && this._emitters[streamName][groupName]) {
   emitter = this._emitters[streamName][groupName];
  }
  else {
   this._client.xgroup('CREATE', streamName, groupName, '0', (err) => {});  //attempt to create Redis group
   emitter = new events.EventEmitter();
   let obj = setInterval(() => {
    this._readGroup(streamName, groupName, consumerName)
    .then(eventList => {
     if (eventList.length > 0) {
      emitter.emit('event', eventList);
     }
    })
    .catch(err => {
     logger.error(`EventStoreClient.subscribe - streamName:${streamName}, groupName:${groupName},\
     consumerName:${consumerName} - ${err}`);
     throw err;     
    });
   }, this._readInterval);
   if (!this._emitters[streamName]) {
    this._emitters[streamName] = {};
   }
   this._emitters[streamName][groupName] = emitter;
   this._intervals.push(obj); 
  }
  
  return emitter;
 }

Sample Results

Below the state of a Redis instance after the following actions:  Create account, Deposit $100, Withdraw $100
127.0.0.1:6379> xrange accountStream - +
1) 1) "1551312621884-0"
   2) 1) "event"
      2) "{\"id\":\"JohnDoe\",\"version\":1,\"type\":\"create\"}"
2) 1) "1551312827949-0"
   2) 1) "event"
      2) "{\"id\":\"JohnDoe\",\"version\":2,\"type\":\"deposit\",\"amount\":100}"
3) 1) "1551312847014-0"
   2) 1) "event"
      2) "{\"id\":\"JohnDoe\",\"version\":3,\"type\":\"withdraw\",\"amount\":100}"
Below is the corresponding state of a MongoDB collection being used for event data aggregation.
> db.accountCollection.find()
{ "_id" : "JohnDoe", "funds" : 0, "timestamps" : [ "1551312621884-0", "1551312827949-0", "1551312847014-0" ] }

Source

Full source w/comments here: https://github.com/joeywhelan/redisStreamEventStore

Copyright ©1993-2024 Joey E Whelan, All rights reserved.

NICE inContact WorkItem Routing


Summary

I'll demonstrate the use of the NICE inContact Work Item API in this post.  That API provides the ability to bring non-native tasks into the routing engine.  The 'workitem' to be routed will be email from Google Gmail.  I'll cover the use of the Gmail API in addition to WorkItem.

This is a contrived example with, at best, demo-grade code.

Overall Architecture

Below is a depiction of the overall set up of this exercise.  Emails are pulled from Gmail via API, submitted to InContact's routing engine via API, then delivered to an agent target via a routing a strategy.


Application Architecture

Below is the application layout.  

Routing Logic

Picture of the routing strategy used for this demo.  This a very simple strategy that pulls a work item from queue and then pops a web page with the email contents to an agent.


Code Snippets


Main Procedure

Code below performs the OAuth2 handshake with Google, gets a list of messages currently in the Inbox, then packages up an array of Promises each of which pulls message contents from Gmail and submit them to the Work Item API.
function processMessage(id) {
 return gmail.getMessage(id)
 .then(msg => {
  return workItem.sendEmail(msg.id, msg.from, msg.payload);
 })
 .then(contactId => {
  return {msgId : id, contactId : contactId};
 });
}

gmail.authorize()
.then(_ => {
 return gmail.getMessageList();
})
.then((msgs) => {
 let p = [];
 msgs.forEach(msg => {
  p.push(processMessage(msg.id));
 });
 return Promise.all(p);
})
.then(results => {
 console.log(results);
})
.catch((err) => {
 console.log(err);
});

Gmail GetMessageList and GetMessage

getMessageList() pulls a list of Gmail IDs currently with the INBOX label.  getMessage() pulls the content of a message for a given Gmail ID.

 getMessageList() {
  console.log('getMessageList()');
  const auth = this.oAuth2Client;
  const gmail = google.gmail({version: 'v1', auth});
  return new Promise((resolve, reject) => {
   gmail.users.messages.list({userId: 'me', labelIds: ['INBOX']}, (err, res) => {
    if (err) {
     reject(err);
    }
    resolve(res);
   });
  })
  .then((res) => {
   return res.data.messages;
  })
  .catch((err) => {
   console.error('getMessageList() - ' + err.message);
   throw err;
  });
 }


 getMessage(id) {
  console.log('getMessage() - id: ' + id);
  const auth = this.oAuth2Client;
  const gmail = google.gmail({version: 'v1', auth});
  return new Promise((resolve, reject) => {
   gmail.users.messages.get({userId: 'me', 'id': id}, (err, res) => {
    if (err) {
     reject(err);
    }
    resolve(res);
   });
  })
  .then((res) => {
   const id = res.data.id;
   const arr = (res.data.payload.headers.find(o => o.name === 'From')).value
   .match(/([a-zA-Z0-9._-]+@[a-zA-Z0-9._-]+\.[a-zA-Z0-9._-]+)/gi);
   let from;
   if (arr) {
    from = arr[0];
   }
   else {
    from = '';
   }
   let payload = '';
   if (res.data.payload.body.data) {
    payload = atob(res.data.payload.body.data);
   }
   return {id: id, from: from, payload: payload};
  })
  .catch((err) => {
   console.error('getMessage() - ' + err.message);
   throw err;
  });
 }


WorkItem API Call

This consists of a simple POST with the workItem params in the body.

 postWorkItem(workItemURL, token, id, from, payload, type) {
  console.log('postWorkItem() - url: ' + workItemURL + ' from: ' + from);
  const body = {
    'pointOfContact': this.poc,
    'workItemId': id,
    'workItemPayload': payload,
    'workItemType': type,
    'from': from
  };
 
  return fetch(workItemURL, {
   method: 'POST',
   body: JSON.stringify(body),
   headers: {
    'Content-Type' : 'application/json', 
    'Authorization' : 'bearer ' + token
   },
   cache: 'no-store',
   mode: 'cors'
  })
  .then(response => {
   if (response.ok) {
    return response.json();
   }
   else {
    const msg = 'response status: ' + response.status;
    throw new Error(msg);
   }
  })
  .then(json => {
    return json.contactId;
  })
  .catch(err => {
   console.error('postWorkItem() - ' + err.message);
   throw err;
  });
 }

Results


Source

Full source w/comments here: https://github.com/joeywhelan/workitem

Copyright ©1993-2024 Joey E Whelan, All rights reserved.

Sunday, February 3, 2019

InContact Custom Report API


Summary

This post will walk through the steps necessary to generate a custom report job on InContact's platform via API.  I'll show examples in Node.js (asynchronous) and Python (synchronous).

Overview of Steps


Get Token

Node.js

 getToken() {
  console.log('getToken()');
  const url = 'https://api.incontact.com/InContactAuthorizationServer/Token';
  const body = {
    'grant_type' : 'password',
    'username' : this.username,
    'password' : this.password
  };
  return fetch(url, {
   method: 'POST',
   body: JSON.stringify(body),
   headers: {
    'Content-Type' : 'application/json', 
    'Authorization' : 'basic ' + this.authCode
   },
   cache: 'no-store',
      mode: 'cors'
  })
  .then(response => {
   if (response.ok) {
    return response.json();
   }
   else {
    const msg = 'response status: ' + response.status;
    throw new Error(msg);
   } 
  })
  .then(json => {
   if (json && json.access_token && json.resource_server_base_uri) {
    return json;
   }
   else {
    const msg = 'missing token and/or uri';
    throw new Error(msg);
   }
  })
  .catch(err => {
   console.error('getToken() - ' + err.message);
   throw err;
  });
 }

Python

    def getToken(self):
        print('getToken()')
        url = 'https://api.incontact.com/InContactAuthorizationServer/Token'
        header = {'Authorization' : b'basic ' + self.authCode, 'Content-Type': 'application/json'}
        body =  {'grant_type' : 'password', 'username' : self.username, 'password' : self.password}
        resp = requests.post(url, headers=header, json=body)
        resp.raise_for_status()
        return resp.json()

Start Report Job

Node.js

 startReportJob(reportId, reportURL, token) {
  const url = reportURL + reportId;
  console.log('startReportJob() - url: ' + url);
  const body = {
    'fileType': 'CSV',
    'includeHeaders': 'true',
    'appendDate': 'true',
    'deleteAfter': '7',
    'overwrite': 'true'
  };
  
  return fetch(url, {
   method: 'POST',
   body: JSON.stringify(body),
   headers: {
    'Content-Type' : 'application/json', 
    'Authorization' : 'bearer ' + token
   },
   cache: 'no-store',
   mode: 'cors'
  })
  .then(response => {
   if (response.ok) {
    return response.json();
   }
   else {
    const msg = 'response status: ' + response.status;
    throw new Error(msg);
   }
  })
  .then(json => {
    return json.jobId;
  })
  .catch(err => {
   console.error('startReportJob() - ' + err.message);
   throw err;
  });
 }

Python

    def startReportJob(self, reportId, reportURL, token):
        url = reportURL + reportId
        print('startReportJob() - url: ' + url);
        body = {
                'fileType': 'CSV',
                'includeHeaders': 'true',
                'appendDate': 'true',
                'deleteAfter': '7',
                'overwrite': 'true'
        }
        header = { 'Content-Type' : 'application/json', 'Authorization' : 'bearer ' + token}
        resp = requests.post(url, headers=header, json=body)
        resp.raise_for_status()
        return resp.json()['jobId']

Get File URL

Node.js

 getFileURL(jobId, reportURL, token, numTries=10) {
  console.log('getFileURL() - jobId: ' + jobId + ' numTries: ' + numTries);
  const that = this;
  const url = reportURL + jobId;
  
  return fetch(url, {
   method: 'GET',
   headers: {
    'Content-Type' : 'application/x-www-form-urlencoded', 
    'Authorization' : 'bearer ' + token
   },
   cache: 'no-store',
   mode: 'cors'
  })
  .then(response => {
   if (response.ok) {
    return response.json();
   }
   else {
    const msg = 'response status: ' + response.status;
    throw new Error(msg);
   }
  })
  .then(json => {
   if (json.jobResult.resultFileURL) {
    return json.jobResult.resultFileURL;
   }
   else {
    if (numTries > 0) {  //loop (recursive) up to the numTries parameter
     return new Promise((resolve, reject) => {
      setTimeout(() => { 
       resolve(that.getFileURL(jobId, reportURL, token, numTries-1));
      }, 60000);  //retry once per minute
     });
    }
    else {
     throw new Error('Maximum retries reached');
    } 
   }
  })
  .catch(err => {
   console.error('getFileURL() - ' + err.message);
   throw err;
  });
 }

Python

    def getFileURL(self, jobId, reportURL, token):
        url = reportURL + jobId
        header = { 'Content-Type' : 'application/x-www-form-urlencoded', 'Authorization' : 'bearer ' + token }
        resp = requests.get(url, headers=header)
        fileURL = resp.json()['jobResult']['resultFileURL']
        numTries = 10
        
        while (not fileURL and numTries > 0):
            print('getFileURL() - jobId: ' + jobId + ' numTries: ' + str(numTries))
            time.sleep(60)
            resp = requests.get(url, headers=header)
            fileURL = resp.json()['jobResult']['resultFileURL']
            numTries -= 1
        
        return fileURL

Download Report

Node.js

 
 downloadReport(url, token) {
  console.log('downLoadReport() - url: ' + url);
  
  return fetch(url, {
   method: 'GET',
   headers: {'Authorization' : 'bearer ' + token},
   cache: 'no-store',
   mode: 'cors'
  })
  .then(response => {
   if (response.ok) {
    return response.json();
   }
   else {
    const msg = 'response status: ' + response.status;
    throw new Error(msg);
   }
  })
  .then(json => {
    return json.files.file;
  })
  .catch(err => {
   console.error('downloadReport() - ' + err.message);
   throw err;
  })

Python

    def downloadReport(self, url, token):
        print('downLoadReport() - url: ' + url)
        header = { 'Content-Type' : 'application/x-www-form-urlencoded', 'Authorization' : 'bearer ' + token }
        resp = requests.get(url, headers=header)
        return resp.json()['files']['file']   

Output

Node.js

getReport() - reportId: 4477
getToken()
startReportJob() - url: https://api-c7.incontact.com/inContactAPI/services/v13.0/report-jobs/4477
getFileURL() - jobId: 825795 numTries: 10
getFileURL() - jobId: 825795 numTries: 9
getFileURL() - jobId: 825795 numTries: 8
getFileURL() - jobId: 825795 numTries: 7
downLoadReport() - url: https://api-C7.incontact.com/inContactAPI/services/V15.0/files?fileName=CustomReports%5cApiReports%5cService+Levels_20190203T054855.csv
Job Complete

Python

getReport() - reportId: 4477
getToken()
startReportJob() - url: https://api-c7.incontact.com/inContactAPI/services/v13.0/report-jobs/4477
getFileURL() - jobId: 825797 numTries: 10
getFileURL() - jobId: 825797 numTries: 9
getFileURL() - jobId: 825797 numTries: 8
downLoadReport() - url: https://api-C7.incontact.com/inContactAPI/services/V15.0/files?fileName=CustomReports%5cApiReports%5cService+Levels_20190203T055156.csv
Job Complete

Source

Full source w/comments - https://github.com/joeywhelan/reportdemo

Copyright ©1993-2024 Joey E Whelan, All rights reserved.