commit 0bf3ed3543b13c400ccaead9bab93836525cf34d
Author: MTRNord <mtrnord1@gmail.com>
Date: Sun, 4 Feb 2018 00:36:01 +0100
Initial Version with very basic Twitter Crawling implemented
Signed-off-by: MTRNord <mtrnord1@gmail.com>
Diffstat:
10 files changed, 268 insertions(+), 0 deletions(-)
diff --git a/.gitignore b/.gitignore
@@ -0,0 +1,70 @@
+# Created by .ignore support plugin (hsz.mobi)
+### JetBrains template
+# Covers JetBrains IDEs: IntelliJ, RubyMine, PhpStorm, AppCode, PyCharm, CLion, Android Studio and Webstorm
+# Reference: https://intellij-support.jetbrains.com/hc/en-us/articles/206544839
+
+# User-specific stuff:
+.idea/**/workspace.xml
+.idea/**/tasks.xml
+.idea/dictionaries
+
+# Sensitive or high-churn files:
+.idea/**/dataSources/
+.idea/**/dataSources.ids
+.idea/**/dataSources.xml
+.idea/**/dataSources.local.xml
+.idea/**/sqlDataSources.xml
+.idea/**/dynamic.xml
+.idea/**/uiDesigner.xml
+
+# Gradle:
+.idea/**/gradle.xml
+.idea/**/libraries
+
+# CMake
+cmake-build-debug/
+cmake-build-release/
+
+# Mongo Explorer plugin:
+.idea/**/mongoSettings.xml
+
+## File-based project format:
+*.iws
+
+## Plugin-specific files:
+
+# IntelliJ
+out/
+
+# mpeltonen/sbt-idea plugin
+.idea_modules/
+
+# JIRA plugin
+atlassian-ide-plugin.xml
+
+# Cursive Clojure plugin
+.idea/replstate.xml
+
+# Crashlytics plugin (for Android Studio and IntelliJ)
+com_crashlytics_export_strings.xml
+crashlytics.properties
+crashlytics-build.properties
+fabric.properties
+### SBT template
+# Simple Build Tool
+# http://www.scala-sbt.org/release/docs/Getting-Started/Directories.html#configuring-version-control
+
+dist/*
+target/
+lib_managed/
+src_managed/
+project/boot/
+project/plugins/project/
+.history
+.cache
+.lib/
+### Scala template
+*.class
+*.log
+src/main/resources/application.conf
+.idea/
+\ No newline at end of file
diff --git a/README.md b/README.md
@@ -0,0 +1 @@
+# FreifunkNews
+\ No newline at end of file
diff --git a/build.sbt b/build.sbt
@@ -0,0 +1,13 @@
+name := "FreifunkNews"
+
+version := "1.0"
+
+scalaVersion := "2.12.3"
+
+resolvers += Resolver.sonatypeRepo("snapshots")
+
+libraryDependencies ++= Seq(
+ "com.danielasfregola" %% "twitter4s" % "5.5-SNAPSHOT",
+ "com.typesafe" % "config" % "1.3.2",
+ "ch.qos.logback" % "logback-classic" % "1.1.9"
+)
diff --git a/project/build.properties b/project/build.properties
@@ -0,0 +1 @@
+sbt.version=1.0.0
diff --git a/project/plugins.sbt b/project/plugins.sbt
diff --git a/src/main/resources/application_example.conf b/src/main/resources/application_example.conf
@@ -0,0 +1,19 @@
+twitter {
+ consumer {
+ key = "my-consumer-key"
+ secret = "my-consumer-secret"
+ }
+ access {
+ key = ""
+ secret = ""
+ }
+ trackedWords = [
+ "#freifunk"
+ ]
+ lists = [
+ "MTRNord/freifunknews"
+ ]
+}
+akka {
+
+}
diff --git a/src/main/scala/Main.scala b/src/main/scala/Main.scala
@@ -0,0 +1,53 @@
+import com.danielasfregola.twitter4s.entities.enums.AccessType
+import com.danielasfregola.twitter4s.entities.enums.AccessType.AccessType
+import com.danielasfregola.twitter4s.entities.{AccessToken, ConsumerToken}
+import com.danielasfregola.twitter4s.util.Configurations.{consumerTokenKey, consumerTokenSecret}
+import com.danielasfregola.twitter4s.{TwitterAuthenticationClient, TwitterRestClient, TwitterStreamingClient}
+
+import scala.concurrent.{Await, Future}
+import scala.concurrent.ExecutionContext.Implicits.global
+import scala.concurrent.duration.DurationInt
+import scala.util.{Failure, Success}
+
+object Main extends App {
+
+ val consumerToken = ConsumerToken(key = consumerTokenKey, secret = consumerTokenSecret)
+ val TwitterAuthClient = new TwitterAuthenticationClient(consumerToken)
+
+ val write_access: AccessType = AccessType.Write
+
+ val reqToken = TwitterAuthClient.requestToken(x_auth_access_type = Some(write_access))
+
+ reqToken onComplete {
+ case Success(token) =>
+ println(TwitterAuthClient.authenticateUrl(token = token.token, force_login = false))
+
+ val pin = scala.io.StdIn.readLine("Insert Pin: ")
+
+ val AccessTokenResp = Await.result(TwitterAuthClient.accessToken(token.token, pin), 10 seconds)
+
+ val consumerToken = ConsumerToken(key = consumerTokenKey, secret = consumerTokenSecret)
+ val accessToken = AccessToken(key = AccessTokenResp.accessToken.key, secret = AccessTokenResp.accessToken.secret)
+
+ val streamingClient = TwitterStreamingClient.apply(accessToken = accessToken, consumerToken = consumerToken)
+ val restClient = TwitterRestClient.apply(accessToken = accessToken, consumerToken = consumerToken)
+
+ val sender = new twitterSender.Sender(restClient = restClient)
+ //sender.sendHelloWorld
+
+ val stream = new twitterCrawler.StreamingApi(streamingClient = streamingClient, restClient = restClient)
+ stream.fetchHastags
+
+ webServer.WebServer.main()
+ case Failure(err) => println(err.toString)
+ }
+
+ //Keep alive workaround
+ val waitFunc = Future {
+ while (true) {
+ Thread.sleep(1000)
+ }
+ }
+
+ Await.result(waitFunc, scala.concurrent.duration.Duration.Inf)
+}
diff --git a/src/main/scala/twitterCrawler/StreamingApi.scala b/src/main/scala/twitterCrawler/StreamingApi.scala
@@ -0,0 +1,52 @@
+package twitterCrawler
+
+import com.danielasfregola.twitter4s.entities.streaming.common.DisconnectMessage
+import com.danielasfregola.twitter4s.entities.{Tweet, User}
+import com.danielasfregola.twitter4s.http.clients.streaming.TwitterStream
+import com.danielasfregola.twitter4s.{TwitterRestClient, TwitterStreamingClient}
+import com.typesafe.config.{Config, ConfigFactory}
+
+import scala.collection.JavaConversions._
+import scala.concurrent.ExecutionContext.Implicits.global
+import scala.concurrent.{Await, Future}
+
+class StreamingApi (val streamingClient: TwitterStreamingClient, val restClient: TwitterRestClient) {
+ val conf: Config = ConfigFactory.load()
+ def fetchHastags: Future[TwitterStream] = {
+ val trackedWords = conf.getStringList("twitter.trackedWords").toSeq
+ val trackedLists = conf.getStringList("twitter.lists").toList
+ val trackedUsers: Seq[Long] = Await.result(this.getListUsers(trackedLists), scala.concurrent.duration.Duration.Inf)
+
+ println(s"Launching streaming session with tracked keywords: $trackedWords\r\n" +
+ s"And with tracked Users: $trackedUsers")
+
+ streamingClient.filterStatuses(tracks = trackedWords, follow = trackedUsers) {
+ case tweet: Tweet => println(tweet.text)
+ case disconnect: DisconnectMessage => println("Disconnect: ", disconnect.disconnect.reason)
+ }
+ }
+
+ private def getListUsers(trackedLists: List[String]): Future[Seq[Long]] = {
+ Future {
+ var trackedUsers: Seq[Long] = Seq()
+ trackedLists.forEach((v: String) => {
+ val splitted = v.split("/")
+ val username = splitted(0)
+ val slug = splitted(1)
+
+ var listUsers = Await.result(restClient.listMembersBySlugAndOwnerName(slug = slug, owner_screen_name = username, include_entities= false), scala.concurrent.duration.Duration.Inf)
+ listUsers.data.users.foreach((user: User) => {
+ trackedUsers = trackedUsers:+user.id
+ })
+ while (listUsers.data.next_cursor != listUsers.data.previous_cursor) {
+ listUsers = Await.result(restClient.listMembersBySlugAndOwnerName(slug = slug, owner_screen_name = username, include_entities= false, cursor = listUsers.data.next_cursor), scala.concurrent.duration.Duration.Inf)
+ listUsers.data.users.foreach((user: User) => {
+ trackedUsers = trackedUsers:+user.id
+ })
+ }
+
+ })
+ trackedUsers
+ }
+ }
+}
+\ No newline at end of file
diff --git a/src/main/scala/twitterSender/Sender.scala b/src/main/scala/twitterSender/Sender.scala
@@ -0,0 +1,11 @@
+package twitterSender
+
+import com.danielasfregola.twitter4s.TwitterRestClient
+
+import scala.concurrent.Await
+
+class Sender (val restClient: TwitterRestClient) {
+ def sendHelloWorld = {
+ Await.result(restClient.createTweet(status = "Hello world vom Scalar Freifunk News Projekt!"), scala.concurrent.duration.Duration.Inf)
+ }
+}
diff --git a/src/main/scala/webServer/WebServer.scala b/src/main/scala/webServer/WebServer.scala
@@ -0,0 +1,45 @@
+package webServer
+
+import akka.actor.ActorSystem
+import akka.http.scaladsl.Http
+import akka.http.scaladsl.model._
+import akka.http.scaladsl.model.HttpMethods._
+import akka.stream.ActorMaterializer
+
+import scala.concurrent.ExecutionContextExecutor
+import scala.io.StdIn
+
+object WebServer {
+ def main() {
+
+ implicit val system: ActorSystem = ActorSystem("FreifunkNews")
+ implicit val materializer: ActorMaterializer = ActorMaterializer()
+ // needed for the future flatMap/onComplete in the end
+ implicit val executionContext: ExecutionContextExecutor = system.dispatcher
+
+ val requestHandler: HttpRequest => HttpResponse = {
+ case HttpRequest(GET, Uri.Path("/"), _, _, _) =>
+ HttpResponse(entity = HttpEntity(
+ ContentTypes.`text/html(UTF-8)`,
+ "<html><body>Hello world!</body></html>"))
+
+ case HttpRequest(GET, Uri.Path("/ping"), _, _, _) =>
+ HttpResponse(entity = "PONG!")
+
+ case HttpRequest(GET, Uri.Path("/crash"), _, _, _) =>
+ sys.error("BOOM!")
+
+ case r: HttpRequest =>
+ r.discardEntityBytes() // important to drain incoming HTTP Entity stream
+ HttpResponse(404, entity = "Unknown resource!")
+ }
+
+ val bindingFuture = Http().bindAndHandleSync(requestHandler, "localhost", 8080)
+
+ println(s"Server online at http://localhost:8080/\nPress RETURN to stop...")
+ StdIn.readLine() // let it run until user presses return
+ bindingFuture
+ .flatMap(_.unbind()) // trigger unbinding from the port
+ .onComplete(_ => system.terminate()) // and shutdown when done
+ }
+}