diff --git a/.gitignore b/.gitignore index adc3638..4b50d5b 100644 --- a/.gitignore +++ b/.gitignore @@ -42,6 +42,7 @@ target/ .cache .project .target/ +.bsp/ # IntelliJ specific .idea/ diff --git a/build.sbt b/build.sbt index 9f380bd..1045fa7 100644 --- a/build.sbt +++ b/build.sbt @@ -16,16 +16,16 @@ val buildSettings = Seq ( exportJars := buildExportJars, updateOptions := updateOptions.value.withCachedResolution(true), shellPrompt := { state => "sbt [%s]> ".format(Project.extract(state).currentProject.id) }, - scalacOptions := Seq("-deprecation", "-unchecked", "-feature", "-target:jvm-1.8", "-language:implicitConversions", "-language:postfixOps", "-Xlint"), - parallelExecution in Test := false, + scalacOptions := Seq("-deprecation", "-unchecked", "-feature", "-release:8", "-language:implicitConversions", "-language:postfixOps", "-Xlint"), + Test / parallelExecution := false, coverageFailOnMinimum := true, coverageOutputHTML := true, - coverageOutputXML := true -) ++ Defaults.itSettings + coverageOutputXML := true, + resolvers += Resolver.mavenCentral +) lazy val project = Project("reactive-sparql", file(".")) - .configs(IntegrationTest) .settings(buildSettings: _*) .settings(name := "reactive-sparql") .settings(libraryDependencies ++= `reactive-sparql-dependencies`) diff --git a/project/Dependencies.scala b/project/Dependencies.scala index bcae342..529a9fd 100644 --- a/project/Dependencies.scala +++ b/project/Dependencies.scala @@ -3,37 +3,42 @@ import sbt._ object Version { - val scala = "2.13.3" - val akka = "2.5.29" - val akkaHttp = "10.1.11" + val scala = "2.13.12" + val pekko = "1.0.2" + val pekkoHttp = "1.0.1" + val sslconfig = "0.6.1" val javaxWsRs = "1.1.1" val rdf4j = "2.1.6" //"2.3.2" val logback = "1.2.3" val scalaTest = "3.0.8" val fuseki = "3.7.0" val xmlBind = "2.3.2" + val xerces = "2.12.2" } object Dependencies { - val akkaActor = "com.typesafe.akka" %% "akka-actor" % Version.akka - val akkaStream = "com.typesafe.akka" %% "akka-stream" % Version.akka - val akkaHttpCore = "com.typesafe.akka" %% "akka-http-core" % Version.akkaHttp - val akkaHttpSprayJson = "com.typesafe.akka" %% "akka-http-spray-json" % Version.akkaHttp - val akkaSlf4j = "com.typesafe.akka" %% "akka-slf4j" % Version.akka + val pekkoActor = "org.apache.pekko" %% "pekko-actor" % Version.pekko + val pekkoStream = "org.apache.pekko" %% "pekko-stream" % Version.pekko + val pekkoHttpCore = "org.apache.pekko" %% "pekko-http-core" % Version.pekkoHttp + val pekkoHttpSprayJson = "org.apache.pekko" %% "pekko-http-spray-json" % Version.pekkoHttp + val pekkoSlf4j = "org.apache.pekko" %% "pekko-slf4j" % Version.pekko + val sslConfigLib = "com.typesafe" %% "ssl-config-core" % Version.sslconfig val javaxWsRs = "javax.ws.rs" % "jsr311-api" % Version.javaxWsRs val logbackClassic = "ch.qos.logback" % "logback-classic" % Version.logback val rdf4jRuntime = "org.eclipse.rdf4j" % "rdf4j-runtime" % Version.rdf4j val jakartaXmlBind = "jakarta.xml.bind" % "jakarta.xml.bind-api" % Version.xmlBind - val scalaTest = "org.scalatest" %% "scalatest" % Version.scalaTest % "it,test" - val akkaTestkit = "com.typesafe.akka" %% "akka-testkit" % Version.akka % "it,test" - val akkaStreamTestkit = "com.typesafe.akka" %% "akka-stream-testkit" % Version.akka % "it,test" - val fusekiServer = "org.apache.jena" % "jena-fuseki-server" % Version.fuseki % "it,test" + val xercesImpl = "xerces" % "xercesImpl" % Version.xerces + + val scalaTest = "org.scalatest" %% "scalatest" % Version.scalaTest % Test + val pekkoTestkit = "org.apache.pekko" %% "pekko-testkit" % Version.pekko % Test + val pekkoStreamTestkit = "org.apache.pekko" %% "pekko-stream-testkit" % Version.pekko % Test + val fusekiServer = "org.apache.jena" % "jena-fuseki-server" % Version.fuseki % Test val `reactive-sparql-dependencies` = Seq( - akkaActor, akkaStream, akkaHttpCore, akkaHttpSprayJson, akkaSlf4j, + pekkoActor, pekkoStream, pekkoHttpCore, pekkoHttpSprayJson, pekkoSlf4j, javaxWsRs, rdf4jRuntime, - logbackClassic, scalaTest, akkaTestkit, akkaStreamTestkit, fusekiServer, - jakartaXmlBind) + logbackClassic, scalaTest, pekkoTestkit, pekkoStreamTestkit, fusekiServer, + jakartaXmlBind, xercesImpl) } diff --git a/project/build.properties b/project/build.properties index c0bab04..04267b1 100644 --- a/project/build.properties +++ b/project/build.properties @@ -1 +1 @@ -sbt.version=1.2.8 +sbt.version=1.9.9 diff --git a/publish.sbt b/publish.sbt index 02a7f03..11d0eac 100644 --- a/publish.sbt +++ b/publish.sbt @@ -40,5 +40,3 @@ pomExtra := { } - -pgpReadOnly := true diff --git a/src/main/scala/ai/agnos/sparql/api/ClientAPIProtocol.scala b/src/main/scala/ai/agnos/sparql/api/ClientAPIProtocol.scala index 98ea7b5..b4947fd 100644 --- a/src/main/scala/ai/agnos/sparql/api/ClientAPIProtocol.scala +++ b/src/main/scala/ai/agnos/sparql/api/ClientAPIProtocol.scala @@ -1,6 +1,6 @@ package ai.agnos.sparql.api -import akka.http.scaladsl.model.HttpMethod +import org.apache.pekko.http.scaladsl.model.HttpMethod /** diff --git a/src/main/scala/ai/agnos/sparql/api/GraphStoreProtocol.scala b/src/main/scala/ai/agnos/sparql/api/GraphStoreProtocol.scala index c097cb1..1b7c4fb 100644 --- a/src/main/scala/ai/agnos/sparql/api/GraphStoreProtocol.scala +++ b/src/main/scala/ai/agnos/sparql/api/GraphStoreProtocol.scala @@ -3,8 +3,8 @@ package ai.agnos.sparql.api import java.net.URL import java.nio.file.Path -import akka.http.scaladsl.model.HttpMethod -import akka.http.scaladsl.model.HttpMethods._ +import org.apache.pekko.http.scaladsl.model.HttpMethod +import org.apache.pekko.http.scaladsl.model.HttpMethods._ import org.eclipse.rdf4j.model.{IRI, Model} import org.eclipse.rdf4j.rio.RDFFormat import org.eclipse.rdf4j.rio.RDFFormat.NTRIPLES diff --git a/src/main/scala/ai/agnos/sparql/api/PrefixMapping.scala b/src/main/scala/ai/agnos/sparql/api/PrefixMapping.scala index 5559650..9c8ba34 100644 --- a/src/main/scala/ai/agnos/sparql/api/PrefixMapping.scala +++ b/src/main/scala/ai/agnos/sparql/api/PrefixMapping.scala @@ -1,6 +1,6 @@ package ai.agnos.sparql.api -import com.sun.org.apache.xerces.internal.util.XMLChar +import org.apache.xerces.util.XMLChar import scala.util.control.Breaks._ @@ -22,7 +22,6 @@ object NamespaceConstants { val RDF = "http://www.w3.org/1999/02/22-rdf-syntax-ns#" val RDFS = "http://www.w3.org/2000/01/rdf-schema#" - val RDFSyntax = "http://www.w3.org/TR/rdf-syntax-grammar#" val OWL = "http://www.w3.org/2002/07/owl#" val DC_11 = "http://purl.org/dc/elements/1.1/" val TERMS = "http://purl.org/dc/terms/" @@ -37,7 +36,7 @@ object PrefixMapping { import NamespaceConstants._ - class IllegalPrefixException(prefix : String) extends IllegalArgumentException + private class IllegalPrefixException(prefix : String) extends IllegalArgumentException(prefix) def none = new PrefixMapping @@ -45,7 +44,7 @@ object PrefixMapping { * A PrefixMapping that contains the "standard" prefixes we know about, * viz rdf, rdfs, dc, rss, vcard, and owl. */ - def standard = { + def standard: PrefixMapping = { val pm = new PrefixMapping pm.setNsPrefix(PREFIX_RDFS, RDFS) pm.setNsPrefix(PREFIX_RDF, RDF) @@ -55,7 +54,7 @@ object PrefixMapping { pm } - def extended = { + def extended: PrefixMapping = { val pm = standard pm.setNsPrefix(PREFIX_SKOS, SKOS) pm.setNsPrefix(PREFIX_FOAF, FOAF) @@ -86,7 +85,7 @@ object PrefixMapping { * @param uri * @return the index of the first character of the localname */ - def splitNamespace(uri : String) : Int = { + private def splitNamespace(uri : String) : Int = { // XML Namespaces 1.0: // A qname name is NCName ':' NCName @@ -115,7 +114,7 @@ object PrefixMapping { while (i >= 1) { i -= 1 ch = uri.charAt(i) - if (notNameChar(ch)) break + if (notNameChar(ch)) break() } var j = i + 1 @@ -152,7 +151,7 @@ object PrefixMapping { // Do a quick test before calling .startsWith // OLD: if ( uri.charAt(j - 1) == ':' && uri.lastIndexOf(':', j - 2) == -1) // - if (!(j == 7 && uri.startsWith("mailto:"))) break + if (!(j == 7 && uri.startsWith("mailto:"))) break() } } j @@ -161,11 +160,11 @@ object PrefixMapping { /** * answer true iff this is not a legal NCName character, ie, is a possible split-point start. */ - def notNameChar(ch : Char) : Boolean = !XMLChar.isNCName(ch) + private def notNameChar(ch : Char) : Boolean = !XMLChar.isNCName(ch) } /** - * Inspired by Jena's PrefixMappingImpl class, this class does more or les the same. + * Inspired by Jena's PrefixMappingImpl class, this class does more or less the same. * * See http://svn.apache.org/repos/asf/jena/trunk/jena-core/src/main/java/com/hp/hpl/jena/shared/impl/PrefixMappingImpl.java */ @@ -173,10 +172,10 @@ class PrefixMapping { import PrefixMapping._ - protected var prefixToURI : Map[String, String] = Map.empty - protected var URItoPrefix : Map[String, String] = Map.empty + private var prefixToURI : Map[String, String] = Map.empty + private var URItoPrefix : Map[String, String] = Map.empty - protected def set(prefix : String, uri : String) { + private def set(prefix : String, uri : String): Unit = { prefixToURI += prefix -> uri URItoPrefix += uri -> prefix } @@ -198,7 +197,7 @@ class PrefixMapping { this } - protected def regenerateReverseMapping() { + private def regenerateReverseMapping(): Unit = { URItoPrefix = prefixToURI.map(_.swap) } @@ -209,7 +208,7 @@ class PrefixMapping { */ def withDefaultMappings(other : PrefixMapping) : PrefixMapping = { - for ((prefix, uri) ← other.prefixToURI) { + for ((prefix, uri) <- other.prefixToURI) { if (getNsPrefixURI(prefix) == null && getNsURIPrefix(uri) == null) { setNsPrefix(prefix, uri) } @@ -226,9 +225,9 @@ class PrefixMapping { * * @param other the Map whose bindings we are to add to this. */ - def setNsPrefixes(other : Map[String, String]) : PrefixMapping = { + private def setNsPrefixes(other : Map[String, String]) : PrefixMapping = { - for ((prefix, uri) ← other) { + for ((prefix, uri) <- other) { setNsPrefix(prefix, uri) } @@ -245,7 +244,7 @@ class PrefixMapping { /** * Checks that a prefix is "legal" - it must be a valid XML NCName. */ - private def checkLegal(prefix : String) { + private def checkLegal(prefix : String): Unit = { if (prefix.length > 0 && !XMLChar.isValidNCName(prefix)) throw new PrefixMapping.IllegalPrefixException(prefix) } @@ -309,7 +308,7 @@ class PrefixMapping { * Answer the qname for uri which uses a prefix from this mapping, or null if there isn't one. *

* Relies on splitNamespace to carve uri into namespace and - * localname components; this ensures that the localname is legal and we just + * localname components; this ensures that the localname is legal, and we just * have to (reverse-)lookup the namespace in the prefix table. *

* @see com.hp.hpl.jena.shared.PrefixMapping#qnameFor(java.lang.String) @@ -330,52 +329,18 @@ class PrefixMapping { null } else { - prefix + ":" + local + s"$prefix:$local" } } /** - * Compress the URI using the prefix mapping. This version of the code looks through all the maplets and checks each - * candidate prefix URI for being a leading substring of the argument URI. There's probably a much more efficient - * algorithm available, pre-processing the prefix strings into some kind of search table, but for the moment we don't - * need it. + * Compress the URI using the prefix mapping. This version of the code looks through all the maplet */ - def shortForm(uri : String) : String = { - val (prefix, otherUri) = findMapping(uri, true) - if (prefix == null) { - uri - } - else { - s"${prefix}:${uri.substring(otherUri.length)}" - } - } - - def samePrefixMappingAs(other : PrefixMapping) : Boolean = prefixToURI == other.prefixToURI - - /** - * Answer a prefixToURI entry in which the value is an initial substring of uri. - * If partial is false, then the value must equal uri. - * - * Does a linear search of the entire prefixToURI, so not terribly efficient for large maps. - * - * @param uri the value to search for - * @param partial true if the match can be any leading substring, false for exact match - * @return some entry (k, v) such that uri starts with v [equal for partial=false] - */ - private def findMapping(uri : String, partial : Boolean) : (String, String) = { - for ((prefix, otherUri) ← prefixToURI) { - if (uri.startsWith(otherUri) && (partial || otherUri.length == uri.length)) { - return (prefix, otherUri) - } - } - (null, null) + def sparql : String = { + pairs.map { + case (key, value) => s"PREFIX ${key}: <${value}>" + }.mkString("\n", "\n", "\n") } private def pairs = prefixToURI.toList sortBy { _._1 } - - def sparql : String = { - pairs map { - case (key, value) ⇒ s"PREFIX ${key}: <${value}>" - } mkString ("\n", "\n", "\n") - } } diff --git a/src/main/scala/ai/agnos/sparql/api/SparqlClientProtocol.scala b/src/main/scala/ai/agnos/sparql/api/SparqlClientProtocol.scala index 7e600a3..06227ae 100644 --- a/src/main/scala/ai/agnos/sparql/api/SparqlClientProtocol.scala +++ b/src/main/scala/ai/agnos/sparql/api/SparqlClientProtocol.scala @@ -1,8 +1,8 @@ package ai.agnos.sparql.api -import akka.actor.ActorSystem -import akka.http.scaladsl.model.{StatusCode, StatusCodes} -import akka.http.scaladsl.model.HttpMethods._ +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.model.{StatusCode, StatusCodes} +import org.apache.pekko.http.scaladsl.model.HttpMethods._ trait SparqlClientProtocol extends ClientAPIProtocol diff --git a/src/main/scala/ai/agnos/sparql/api/SparqlConstruct.scala b/src/main/scala/ai/agnos/sparql/api/SparqlConstruct.scala index 02fe663..e4151d2 100644 --- a/src/main/scala/ai/agnos/sparql/api/SparqlConstruct.scala +++ b/src/main/scala/ai/agnos/sparql/api/SparqlConstruct.scala @@ -1,7 +1,7 @@ package ai.agnos.sparql.api -import akka.http.scaladsl.model.{HttpMethod, HttpMethods} +import org.apache.pekko.http.scaladsl.model.{HttpMethod, HttpMethods} import org.eclipse.rdf4j.model.{BNode, IRI, Literal, Value} @@ -120,7 +120,8 @@ abstract class SparqlConstruct()( case literal: Literal => val tpe = literal.getDatatype s"'${literal.stringValue()}'^^${pm.getNsURIPrefix(tpe.getNamespace)}:${tpe.getLocalName}" - case bn: BNode => throw new IllegalArgumentException("Should not use Blank Node as query parameter") + case _: BNode => throw new IllegalArgumentException("Should not use Blank Node as query parameter") + case _ => throw new IllegalArgumentException(s"Unsupported value type ${value.getClass.getName}") } } } diff --git a/src/main/scala/ai/agnos/sparql/api/SparqlQuery.scala b/src/main/scala/ai/agnos/sparql/api/SparqlQuery.scala index 5919b04..003b306 100644 --- a/src/main/scala/ai/agnos/sparql/api/SparqlQuery.scala +++ b/src/main/scala/ai/agnos/sparql/api/SparqlQuery.scala @@ -1,8 +1,8 @@ package ai.agnos.sparql.api -import akka.http.scaladsl.model._ -import akka.stream.scaladsl.Source -import akka.util.ByteString +import org.apache.pekko.http.scaladsl.model._ +import org.apache.pekko.stream.scaladsl.Source +import org.apache.pekko.util.ByteString import org.eclipse.rdf4j.model.{IRI, Value} import ai.agnos.sparql.stream.client.SparqlClientConstants._ import ai.agnos.sparql.util.SparqlQueryStringConverter diff --git a/src/main/scala/ai/agnos/sparql/api/SparqlStatement.scala b/src/main/scala/ai/agnos/sparql/api/SparqlStatement.scala index bcde9b5..fd680ca 100644 --- a/src/main/scala/ai/agnos/sparql/api/SparqlStatement.scala +++ b/src/main/scala/ai/agnos/sparql/api/SparqlStatement.scala @@ -1,9 +1,8 @@ package ai.agnos.sparql.api -import akka.http.scaladsl.model._ +import org.apache.pekko.http.scaladsl.model._ import scala.concurrent.duration._ -import scala.language.postfixOps /** diff --git a/src/main/scala/ai/agnos/sparql/api/SparqlUpdate.scala b/src/main/scala/ai/agnos/sparql/api/SparqlUpdate.scala index e4d0836..4a92374 100644 --- a/src/main/scala/ai/agnos/sparql/api/SparqlUpdate.scala +++ b/src/main/scala/ai/agnos/sparql/api/SparqlUpdate.scala @@ -2,7 +2,7 @@ package ai.agnos.sparql.api import java.text.SimpleDateFormat -import akka.http.scaladsl.model.{HttpMethod, HttpMethods} +import org.apache.pekko.http.scaladsl.model.{HttpMethod, HttpMethods} object SparqlUpdate { diff --git a/src/main/scala/ai/agnos/sparql/mapper/DelegatingSolutionMapper.scala b/src/main/scala/ai/agnos/sparql/mapper/DelegatingSolutionMapper.scala index aa0ff95..632cd58 100644 --- a/src/main/scala/ai/agnos/sparql/mapper/DelegatingSolutionMapper.scala +++ b/src/main/scala/ai/agnos/sparql/mapper/DelegatingSolutionMapper.scala @@ -6,7 +6,7 @@ import ai.agnos.sparql.api.QuerySolution /** * A helper mapper that delegates mapping to a specified function. */ -class DelegatingSolutionMapper[T] private (mapper : QuerySolution ⇒ T) +class DelegatingSolutionMapper[T] private (mapper : QuerySolution => T) extends SolutionMapper[T] { def map(querySolution : QuerySolution) : T = { @@ -15,5 +15,5 @@ class DelegatingSolutionMapper[T] private (mapper : QuerySolution ⇒ T) } object DelegatingSolutionMapper { - def apply[T](mapper : QuerySolution ⇒ T) = new DelegatingSolutionMapper[T](mapper) + def apply[T](mapper : QuerySolution => T) = new DelegatingSolutionMapper[T](mapper) } diff --git a/src/main/scala/ai/agnos/sparql/mapper/SparqlClientJsonProtocol.scala b/src/main/scala/ai/agnos/sparql/mapper/SparqlClientJsonProtocol.scala index ee74082..842f46e 100644 --- a/src/main/scala/ai/agnos/sparql/mapper/SparqlClientJsonProtocol.scala +++ b/src/main/scala/ai/agnos/sparql/mapper/SparqlClientJsonProtocol.scala @@ -1,8 +1,8 @@ package ai.agnos.sparql.mapper -import akka.http.scaladsl.unmarshalling._ +import org.apache.pekko.http.scaladsl.unmarshalling._ import ai.agnos.sparql.api._ -import akka.http.scaladsl.marshallers.sprayjson.SprayJsonSupport +import org.apache.pekko.http.scaladsl.marshallers.sprayjson.SprayJsonSupport import spray.json._ /** @@ -14,21 +14,21 @@ object SparqlClientJsonProtocol extends SprayJsonSupport with DefaultJsonProtoco * A set if JSON parsers form "query results" format compliant * with https://www.w3.org/TR/2013/REC-sparql11-results-json-20130321/ */ - implicit val format4 = jsonFormat3(QuerySolutionValue) + implicit val format4: RootJsonFormat[QuerySolutionValue] = jsonFormat3(QuerySolutionValue) implicit object format5 extends RootJsonFormat[QuerySolution] { def write(c : QuerySolution) = { JsObject(c.values.map(e => e._1 -> e._2.toJson)) } def read(row : JsValue) = read(row.asInstanceOf[JsObject]) def read(row : JsObject) = { - QuerySolution((row.fields mapValues { + QuerySolution((row.fields.view.mapValues { (value : JsValue) => value.convertTo[QuerySolutionValue] - }).toMap) + }.toMap)) } } - implicit val format2 = jsonFormat(ResultSetResults, "bindings") - implicit val format1 = jsonFormat(ResultSetVars, "vars") - implicit val format3 = jsonFormat2(ResultSet) + implicit val format2: RootJsonFormat[ResultSetResults] = jsonFormat(ResultSetResults, "bindings") + implicit val format1: RootJsonFormat[ResultSetVars] = jsonFormat(ResultSetVars, "vars") + implicit val format3: RootJsonFormat[ResultSet] = jsonFormat2(ResultSet) /** @@ -36,7 +36,7 @@ object SparqlClientJsonProtocol extends SprayJsonSupport with DefaultJsonProtoco * @return */ implicit def booleanEntityUnmarshaller: FromEntityUnmarshaller[Boolean] = - byteStringUnmarshaller mapWithInput { (entity, bytes) ⇒ + byteStringUnmarshaller mapWithInput { (entity, bytes) => if (entity.isKnownEmpty) false else bytes.decodeString(Unmarshaller.bestUnmarshallingCharsetFor(entity).nioCharset.name).toLowerCase.equals("true") } diff --git a/src/main/scala/ai/agnos/sparql/stream/client/GraphStoreRequestFlowBuilder.scala b/src/main/scala/ai/agnos/sparql/stream/client/GraphStoreRequestFlowBuilder.scala index 4eac8d1..9fff15e 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/GraphStoreRequestFlowBuilder.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/GraphStoreRequestFlowBuilder.scala @@ -6,14 +6,14 @@ import java.io.{StringReader, StringWriter} import java.net.URL import java.nio.file.Path -import akka.NotUsed -import akka.actor.ActorSystem -import akka.http.scaladsl.Http -import akka.http.scaladsl.model.headers.Accept -import akka.http.scaladsl.model.{HttpEntity, _} -import akka.stream.ActorMaterializer -import akka.stream.scaladsl.{FileIO, Flow, Framing, Source} -import akka.util.ByteString +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.model.headers.Accept +import org.apache.pekko.http.scaladsl.model.{HttpEntity, _} +import org.apache.pekko.stream.Materializer +import org.apache.pekko.stream.scaladsl.{FileIO, Flow, Framing, Source} +import org.apache.pekko.util.ByteString import ai.agnos.sparql.api._ import ai.agnos.sparql.util.HttpEndpoint import org.eclipse.rdf4j.model.{IRI, Model} @@ -67,7 +67,7 @@ trait GraphStoreRequestFlowBuilder extends SparqlClientHelpers with HttpClientFl import GraphStoreRequestFlowBuilder._ implicit val system: ActorSystem - implicit val materializer: ActorMaterializer + implicit val materializer: Materializer /** @@ -114,7 +114,11 @@ trait GraphStoreRequestFlowBuilder extends SparqlClientHelpers with HttpClientFl } def makeModelSource(entity: HttpEntity): Source[Option[Model], Any] = { - if ( !entity.isChunked() && (entity.isKnownEmpty() || entity.contentLengthOption.getOrElse(0) == 0)) { + if ( !entity.isChunked() && + ( + entity.isKnownEmpty() + || entity.contentLengthOption.map(_.toLong).getOrElse(0L) == 0) + ) { // if we know there are no bytes in the entity (no-graph has been returned) // or the reponse content type is not what we have requested then no model is emitted. entity.discardBytes() @@ -198,6 +202,9 @@ trait GraphStoreRequestFlowBuilder extends SparqlClientHelpers with HttpClientFl makeInsertGraphHttpRequest(endpoint, method, graphUri, mapRdfFormatToContentType(format)) { () => makeGraphSource(path, format) } + + case _ => + throw new IllegalArgumentException(s"unsupported graph store request: ${request}") } } diff --git a/src/main/scala/ai/agnos/sparql/stream/client/HttpClientFlowBuilder.scala b/src/main/scala/ai/agnos/sparql/stream/client/HttpClientFlowBuilder.scala index a6b4e57..263510e 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/HttpClientFlowBuilder.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/HttpClientFlowBuilder.scala @@ -1,14 +1,14 @@ package ai.agnos.sparql.stream.client -import akka.NotUsed -import akka.actor.ActorSystem -import akka.http.scaladsl.model.{HttpRequest, HttpResponse} -import akka.http.scaladsl.{Http, HttpsConnectionContext} -import akka.stream.{ActorMaterializer, OverflowStrategy, QueueOfferResult } -import akka.stream.scaladsl.{Flow, Keep, Sink, Source} +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.model.{HttpRequest, HttpResponse} +import org.apache.pekko.http.scaladsl.{ConnectionContext, Http, HttpsConnectionContext} +import org.apache.pekko.stream.{Materializer, OverflowStrategy, QueueOfferResult } +import org.apache.pekko.stream.scaladsl.{Flow, Keep, Sink, Source} import ai.agnos.sparql.util.HttpEndpoint -import com.typesafe.sslconfig.akka.AkkaSSLConfig +import javax.net.ssl.SSLContext import scala.concurrent.{Future, Promise} import scala.util.{Failure, Success, Try} @@ -16,23 +16,26 @@ import scala.util.{Failure, Success, Try} trait HttpClientFlowBuilder { def defaultHttpClientFlow[T](endpoint: HttpEndpoint) - (implicit system: ActorSystem, materializer: ActorMaterializer) + (implicit system: ActorSystem, materializer: Materializer) : Flow[(HttpRequest, T), (Try[HttpResponse], T), NotUsed] = { queuedAndPooledHttpClientFlow[T](endpoint, queueSize = 10, overflowStrategy = OverflowStrategy.backpressure) } /** - * Returns the default - * @param endpoint - * @tparam T - * @return + * Returns the default pooled HTTP client flow + * @param endpoint the endpoint to connect to + * @param httpsContext optional HTTPS connection context + * @tparam T the type parameter + * @return the flow */ - def pooledHttpClientFlow[T](endpoint: HttpEndpoint, akkaSSLConfig: Option[AkkaSSLConfig] = None) - (implicit system: ActorSystem, materializer: ActorMaterializer) + def pooledHttpClientFlow[T](endpoint: HttpEndpoint, httpsContext: Option[HttpsConnectionContext] = None) + (implicit system: ActorSystem) : Flow[(HttpRequest, T), (Try[HttpResponse], T), NotUsed] = { endpoint.protocol match { case "http" => Flow[(HttpRequest, T)].via(Http().cachedHostConnectionPool(endpoint.host, endpoint.port)) - case "https" => Flow[(HttpRequest, T)].via(Http().cachedHostConnectionPoolHttps(endpoint.host, endpoint.port, httpsConnectionContext(akkaSSLConfig))) + case "https" => + val ctx = httpsContext.getOrElse(ConnectionContext.httpsClient(SSLContext.getDefault)) + Flow[(HttpRequest, T)].via(Http().cachedHostConnectionPoolHttps(endpoint.host, endpoint.port, ctx)) case protocol => throw new IllegalArgumentException(s"invalid protocol specified: ${protocol}") } } @@ -45,28 +48,28 @@ trait HttpClientFlowBuilder { * @param endpoint the HTTP(S) endpoint * @param queueSize the size of the queue * @param overflowStrategy the overflow strategy to apply if the queue overruns - * @param akkaSSLConfig SSL config - * @param system - * @param materializer - * @tparam T - * @return + * @param httpsContext optional HTTPS connection context + * @param system the actor system + * @param materializer the materializer + * @tparam T the type parameter + * @return the flow */ def queuedAndPooledHttpClientFlow[T] ( endpoint: HttpEndpoint, queueSize: Int = 10, overflowStrategy: OverflowStrategy = OverflowStrategy.backpressure, - akkaSSLConfig: Option[AkkaSSLConfig] = None + httpsContext: Option[HttpsConnectionContext] = None ) ( - implicit system: ActorSystem, materializer: ActorMaterializer + implicit system: ActorSystem, materializer: Materializer ): Flow[(HttpRequest, T), (Try[HttpResponse], T), NotUsed] = { implicit val _ec = system.dispatcher // Materialize a queue with the desired buffer and overflow strategy val queue = Source.queue[(HttpRequest, Promise[HttpResponse])](queueSize, overflowStrategy) - .via(pooledHttpClientFlow[Promise[HttpResponse]](endpoint, akkaSSLConfig)) + .via(pooledHttpClientFlow[Promise[HttpResponse]](endpoint, httpsContext)) .toMat(Sink.foreach({ case ((Success(resp), p)) => p.success(resp) case ((Failure(e), p)) => p.failure(e) @@ -100,27 +103,4 @@ trait HttpClientFlowBuilder { } flow } - - /** - * Override if the default connection context is not satisfactory. - * - * @return - */ - protected def httpsConnectionContext - ( - config: Option[AkkaSSLConfig] - ) - (implicit system: ActorSystem, materializer: ActorMaterializer): HttpsConnectionContext = { - Http().createClientHttpsContext(config.getOrElse(defaultSSLConfig)) - } - - /** - * The default HTTPS connection context config is very slack, use only for tests - not very useful for real world use. - */ - protected def defaultSSLConfig() - (implicit system: ActorSystem, materializer: ActorMaterializer) - : AkkaSSLConfig = { - AkkaSSLConfig().mapSettings(s => s.withLoose(s.loose.withDisableSNI(true))) - } - } diff --git a/src/main/scala/ai/agnos/sparql/stream/client/HttpEndpointFlow.scala b/src/main/scala/ai/agnos/sparql/stream/client/HttpEndpointFlow.scala index e1441b7..6c6ec4d 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/HttpEndpointFlow.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/HttpEndpointFlow.scala @@ -1,8 +1,8 @@ package ai.agnos.sparql.stream.client -import akka.NotUsed -import akka.http.scaladsl.model.{HttpRequest, HttpResponse} -import akka.stream.scaladsl.Flow +import org.apache.pekko.NotUsed +import org.apache.pekko.http.scaladsl.model.{HttpRequest, HttpResponse} +import org.apache.pekko.stream.scaladsl.Flow import ai.agnos.sparql.util.HttpEndpoint import scala.util.Try diff --git a/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientConstants.scala b/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientConstants.scala index 131582d..493e52d 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientConstants.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientConstants.scala @@ -1,7 +1,7 @@ package ai.agnos.sparql.stream.client -import akka.http.scaladsl.model.MediaType.NotCompressible -import akka.http.scaladsl.model.{ContentType, HttpCharsets, MediaType} +import org.apache.pekko.http.scaladsl.model.MediaType.NotCompressible +import org.apache.pekko.http.scaladsl.model.{ContentType, HttpCharsets, MediaType} import org.eclipse.rdf4j.model.{ModelFactory, ValueFactory} import org.eclipse.rdf4j.model.impl.{LinkedHashModelFactory, SimpleValueFactory} diff --git a/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientHelpers.scala b/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientHelpers.scala index 83d22fc..8aefebd 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientHelpers.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/SparqlClientHelpers.scala @@ -2,11 +2,11 @@ package ai.agnos.sparql.stream.client import ai.agnos.sparql._ -import akka.actor.ActorSystem -import akka.http.scaladsl.model.{HttpEntity, _} -import akka.http.scaladsl.model.headers.{Accept, Authorization, BasicHttpCredentials} -import akka.http.scaladsl.unmarshalling.{FromEntityUnmarshaller, PredefinedFromEntityUnmarshallers} -import akka.stream.ActorMaterializer +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.model.{HttpEntity, _} +import org.apache.pekko.http.scaladsl.model.headers.{Accept, Authorization, BasicHttpCredentials} +import org.apache.pekko.http.scaladsl.unmarshalling.{FromEntityUnmarshaller, PredefinedFromEntityUnmarshallers} +import org.apache.pekko.stream.Materializer import ai.agnos.sparql.api._ import ai.agnos.sparql.stream.client.SparqlClientConstants._ import ai.agnos.sparql.util.{BasicAuthentication, HttpEndpoint} @@ -20,7 +20,7 @@ trait SparqlClientHelpers { import HttpMethods._ implicit val system: ActorSystem - implicit val materializer: ActorMaterializer + implicit val materializer: Materializer implicit val rawBooleanFromEntityUnmarshaller: FromEntityUnmarshaller[Boolean] = PredefinedFromEntityUnmarshallers.stringUnmarshaller.map(_.toBoolean) @@ -32,7 +32,7 @@ trait SparqlClientHelpers { def acceptQueryMediaType(queryType: QueryType): MediaType = { val ct = queryType match { case StreamedQuery(contentType) => contentType - case _: MappedQuery[_] => `application/sparql-results+json` + case _ => `application/sparql-results+json` } ct.mediaType } @@ -65,7 +65,6 @@ trait SparqlClientHelpers { s"$QUERY_PARAM_NAME=${urlEncode(query)}&$REASONING_PARAM_NAME=$reasoning" ) - case SparqlUpdate(POST, update) => HttpRequest( method = HttpMethods.POST, @@ -75,6 +74,9 @@ trait SparqlClientHelpers { `application/x-www-form-urlencoded`, s"$UPDATE_PARAM_NAME=${urlEncode(update)}" ) + + case _ => + throw new IllegalArgumentException(s"unsupported SPARQL statement: ${statement}") } // JC: the method name could be more specific @@ -108,6 +110,9 @@ trait SparqlClientHelpers { case (Failure(throwable), request) => val error = SparqlClientRequestFailedWithError("Request failed on the HTTP layer", throwable) SparqlResponse(success = false, request = request, error = Some(error)) + case _ => + val error = SparqlClientRequestFailed("Request failed with unknown error") + SparqlResponse(success = false, request = response._2, error = Some(error)) } def mapRdfFormatToContentType(format: RDFFormat): ContentType = format match { @@ -119,6 +124,8 @@ trait SparqlClientHelpers { `text/turtle` case f: RDFFormat if f == RDFFormat.JSONLD => `application/ld+json` + case _ => + throw new IllegalArgumentException(s"unsupported RDF format: ${format.getName}") } def mapContentTypeToRdfFormat(contentType: ContentType): RDFFormat = { diff --git a/src/main/scala/ai/agnos/sparql/stream/client/SparqlConstructFlowBuilder.scala b/src/main/scala/ai/agnos/sparql/stream/client/SparqlConstructFlowBuilder.scala index b3c2595..c4d9f55 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/SparqlConstructFlowBuilder.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/SparqlConstructFlowBuilder.scala @@ -1,11 +1,11 @@ package ai.agnos.sparql.stream.client -import akka.NotUsed -import akka.actor.ActorSystem -import akka.http.scaladsl.model.{HttpResponse, StatusCodes} -import akka.stream.{ActorMaterializer, FlowShape} -import akka.stream.scaladsl.{Broadcast, Flow, Framing, GraphDSL, Merge, Partition, Source, ZipWith} -import akka.util.ByteString +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.model.{HttpResponse, StatusCodes} +import org.apache.pekko.stream.{Materializer, FlowShape} +import org.apache.pekko.stream.scaladsl.{Broadcast, Flow, Framing, GraphDSL, Merge, Partition, Source, ZipWith} +import org.apache.pekko.util.ByteString import ai.agnos.sparql.api._ import ai.agnos.sparql.stream.client.SparqlClientConstants.{modelFactory => mf, valueFactory => vf} import org.eclipse.rdf4j.model.{IRI, Model, Resource} @@ -27,7 +27,7 @@ trait SparqlConstructFlowBuilder extends SparqlClientHelpers with ErrorHandlerSu import SparqlConstructFlowBuilder._ implicit val system: ActorSystem - implicit val materializer: ActorMaterializer + implicit val materializer: Materializer implicit val dispatcher: ExecutionContext type Sparql = String @@ -57,7 +57,7 @@ trait SparqlConstructFlowBuilder extends SparqlClientHelpers with ErrorHandlerSu } def deReifyConstructSubGraph(): Flow[Model, Model, NotUsed] = { - import scala.collection.JavaConverters._ + import scala.jdk.CollectionConverters._ Flow[Model] .mapConcat(_.stream().iterator().asScala.toList) .sliding(4,4) @@ -81,6 +81,12 @@ trait SparqlConstructFlowBuilder extends SparqlClientHelpers with ErrorHandlerSu // TODO: Add support for sliding through the entity 4 lines at a time (see responseToPagingModelFlow) case (Success(HttpResponse(StatusCodes.OK, _, entity, _)), _) => entity.withoutSizeLimit().dataBytes.fold(ByteString.empty)(_ ++ _).zip(Source.single(entity.contentType)) + case (Failure(err), req) => + errorHandler.handleError(err) + Source.empty + case _ => + errorHandler.handleError(new IllegalArgumentException("unexpected response when processing flow for request")) + Source.empty } .map { x => Rio.parse(x._1.iterator.asInputStream, "", mapContentTypeToRdfFormat(x._2)) @@ -93,25 +99,30 @@ trait SparqlConstructFlowBuilder extends SparqlClientHelpers with ErrorHandlerSu * This flow will also consume the Http response entity if it is received via a valid HTTP response with * an HTTP status code and produces a corresponding SparqlErrorResult */ - lazy val responseToFailureFlow: Flow[(Try[HttpResponse], SparqlRequest), SparqlResult, NotUsed] = { + private lazy val responseToFailureFlow: Flow[(Try[HttpResponse], SparqlRequest), SparqlResult, NotUsed] = { Flow[(Try[HttpResponse], SparqlRequest)] .flatMapConcat { case (Success(HttpResponse(code, _, entity, _)), _) => - entity.withoutSizeLimit().dataBytes.fold(ByteString.empty)(_ ++ _).map { - case message: ByteString => SparqlErrorResult( + entity.withoutSizeLimit().dataBytes.fold(ByteString.empty)(_ ++ _).map { message: ByteString => + SparqlErrorResult( error = new RuntimeException( s"${message.utf8String}" ), code = code.intValue(), - message = "SPARQL endpoint returned unexpected response body") + message = "SPARQL endpoint returned unexpected response body" + ) } case (Failure(err), req) => errorHandler.handleError(err) Source.single(SparqlErrorResult(err, 0, s"unexpected error when processing flow for request: ${req}")) + case _ => + val err = new IllegalArgumentException("unexpected response when processing flow for request") + errorHandler.handleError(err) + Source.single(SparqlErrorResult(err, 0, err.getMessage)) } } - val responseToResultFlow: Flow[(Try[HttpResponse], SparqlRequest), SparqlResult, NotUsed] = { + private val responseToResultFlow: Flow[(Try[HttpResponse], SparqlRequest), SparqlResult, NotUsed] = { Flow.fromGraph(GraphDSL.create() { implicit builder => import GraphDSL.Implicits._ @@ -130,11 +141,14 @@ trait SparqlConstructFlowBuilder extends SparqlClientHelpers with ErrorHandlerSu } - @deprecated + @deprecated("Use responseToSuccessFlow instead", "1.0.0") lazy val responseToPagingModelFlow: Flow[(Try[HttpResponse], SparqlRequest), SparqlResult, NotUsed] = { Flow[(Try[HttpResponse], SparqlRequest)] .flatMapConcat { case (Success(HttpResponse(StatusCodes.OK, _, entity, _)), _) => entity.withoutSizeLimit().getDataBytes() + case _ => + errorHandler.handleError(new IllegalArgumentException("unexpected response when processing flow for request")) + Source.empty } .via(Framing.delimiter(ByteString.fromString("\n"), 1*1024, allowTruncation = true)) .map( bs => Rio.parse(bs.iterator.asInputStream, "", RDFFormat.NQUADS).stream().findFirst().get()) diff --git a/src/main/scala/ai/agnos/sparql/stream/client/SparqlQueryFlowBuilder.scala b/src/main/scala/ai/agnos/sparql/stream/client/SparqlQueryFlowBuilder.scala index f3b17b8..8e1fd21 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/SparqlQueryFlowBuilder.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/SparqlQueryFlowBuilder.scala @@ -1,10 +1,10 @@ package ai.agnos.sparql.stream.client -import akka.actor.ActorSystem -import akka.http.scaladsl.model._ -import akka.stream.{ActorMaterializer, FlowShape} -import akka.stream.scaladsl.{Flow, GraphDSL, Merge, Partition, Source} -import akka.util.ByteString +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.model._ +import org.apache.pekko.stream.{Materializer, FlowShape} +import org.apache.pekko.stream.scaladsl.{Flow, GraphDSL, Merge, Partition, Source} +import org.apache.pekko.util.ByteString import ai.agnos.sparql.api._ import scala.util.{Failure, Success, Try} @@ -18,7 +18,7 @@ trait SparqlQueryFlowBuilder extends SparqlClientHelpers with ErrorHandlerSuppor import ai.agnos.sparql.mapper.SparqlClientJsonProtocol._ implicit val system: ActorSystem - implicit val materializer: ActorMaterializer + implicit val materializer: Materializer implicit val dispatcher: ExecutionContext def sparqlQueryFlow(endpointFlow: HttpEndpointFlow[SparqlRequest]): Flow[SparqlRequest, SparqlResponse, Any] = { @@ -30,6 +30,7 @@ trait SparqlQueryFlowBuilder extends SparqlClientHelpers with ErrorHandlerSuppor val partition = builder.add(Partition[SparqlRequest](routes, { case SparqlRequest(SparqlQuery(_, _, StreamedQuery(_),_,_,_,_,_,_,_), _) => 0 case SparqlRequest(SparqlQuery(_, _,_: MappedQuery[_],_,_,_,_,_,_,_), _) => 1 + case _ => throw new IllegalArgumentException("unsupported query type") })) val responseMerger = builder.add(Merge[SparqlResponse](routes).named("merge.sparqlResponse")) @@ -55,7 +56,7 @@ trait SparqlQueryFlowBuilder extends SparqlClientHelpers with ErrorHandlerSuppor _, List(StreamingSparqlResult(dataStream, Some(contentType))), _ ) if isSparqlResultsJson(contentType) => - Source.fromFuture { + Source.future { dataStream.runFold(ByteString.empty)(_ ++ _).map { data => Try(format3.read(data.utf8String.parseJson)) match { case Success(x: ResultSet) => @@ -98,7 +99,8 @@ trait SparqlQueryFlowBuilder extends SparqlClientHelpers with ErrorHandlerSuppor .map { case request@SparqlRequest(query: SparqlQuery, _) => (makeHttpRequest(endpointFlow.endpoint, query), request) - } + case _ => throw new IllegalArgumentException("unsupported request type") + } .log("SPARQL endpoint request") .via(endpointFlow.flow) .log("SPARQL endpoint response") @@ -107,6 +109,8 @@ trait SparqlQueryFlowBuilder extends SparqlClientHelpers with ErrorHandlerSuppor SparqlResponse(request, status == StatusCodes.OK, status, result = StreamingSparqlResult(entity.dataBytes, Some(entity.contentType)) :: Nil) case (Failure(error), request) => SparqlResponse(request, success = false, error = Some(SparqlClientRequestFailedWithError("failed to execute sparql query", error))) + case _ => + throw new IllegalArgumentException("unexpected response when processing flow for request") } } diff --git a/src/main/scala/ai/agnos/sparql/stream/client/SparqlRequestFlowBuilder.scala b/src/main/scala/ai/agnos/sparql/stream/client/SparqlRequestFlowBuilder.scala index a4e03d8..bccb421 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/SparqlRequestFlowBuilder.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/SparqlRequestFlowBuilder.scala @@ -1,8 +1,8 @@ package ai.agnos.sparql.stream.client -import akka.NotUsed -import akka.stream.FlowShape -import akka.stream.scaladsl.{Flow, GraphDSL, Merge, Partition} +import org.apache.pekko.NotUsed +import org.apache.pekko.stream.FlowShape +import org.apache.pekko.stream.scaladsl.{Flow, GraphDSL, Merge, Partition} import ai.agnos.sparql.api.{SparqlConstruct, _} @@ -26,6 +26,7 @@ trait SparqlRequestFlowBuilder extends SparqlQueryFlowBuilder case SparqlRequest(_: SparqlQuery, _) => 0 case SparqlRequest(_: SparqlUpdate, _) => 1 case SparqlRequest(_: SparqlConstruct, _) => 2 + case _ => throw new IllegalArgumentException("unsupported query type") })) val responseMerger = builder.add(Merge[SparqlResponse](routes).named("merge.sparqlResponse")) diff --git a/src/main/scala/ai/agnos/sparql/stream/client/SparqlUpdateFlowBuilder.scala b/src/main/scala/ai/agnos/sparql/stream/client/SparqlUpdateFlowBuilder.scala index d09571f..60afcd8 100644 --- a/src/main/scala/ai/agnos/sparql/stream/client/SparqlUpdateFlowBuilder.scala +++ b/src/main/scala/ai/agnos/sparql/stream/client/SparqlUpdateFlowBuilder.scala @@ -1,11 +1,11 @@ package ai.agnos.sparql.stream.client -import akka.NotUsed -import akka.actor.ActorSystem -import akka.http.scaladsl.model.{HttpResponse, StatusCodes} -import akka.http.scaladsl.unmarshalling.Unmarshal -import akka.stream.{ActorMaterializer, FlowShape} -import akka.stream.scaladsl.{Broadcast, Flow, GraphDSL, ZipWith} +import org.apache.pekko.NotUsed +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.model.{HttpResponse, StatusCodes} +import org.apache.pekko.http.scaladsl.unmarshalling.Unmarshal +import org.apache.pekko.stream.{Materializer, FlowShape} +import org.apache.pekko.stream.scaladsl.{Broadcast, Flow, GraphDSL, ZipWith} import ai.agnos.sparql.api._ import scala.concurrent.{ExecutionContext, Future} @@ -16,7 +16,7 @@ trait SparqlUpdateFlowBuilder extends SparqlClientHelpers with ErrorHandlerSuppo import SparqlClientConstants._ implicit val system: ActorSystem - implicit val materializer: ActorMaterializer + implicit val materializer: Materializer implicit val dispatcher: ExecutionContext def sparqlUpdateFlow(endpointFlow: HttpEndpointFlow[SparqlRequest]): Flow[SparqlRequest, SparqlResponse, NotUsed] = { diff --git a/src/main/scala/ai/agnos/sparql/util/SparqlQueryStringConverter.scala b/src/main/scala/ai/agnos/sparql/util/SparqlQueryStringConverter.scala index fbb9b32..0106ca6 100644 --- a/src/main/scala/ai/agnos/sparql/util/SparqlQueryStringConverter.scala +++ b/src/main/scala/ai/agnos/sparql/util/SparqlQueryStringConverter.scala @@ -23,6 +23,7 @@ object SparqlQueryStringConverter { urlEncode(l.toString) case i: IRI => urlEncode(s"<${i.toString}>") case _: BNode => throw new IllegalArgumentException("BNode bindings are not allowed") + case _ => throw new IllegalArgumentException(s"unsupported binding value type [${value.getClass.getName}]") } } def mkColBindings[T](bindVar: String, value: Iterable[T]): Iterable[String] = { diff --git a/src/test/resources/application.conf b/src/test/resources/application.conf index 9184ad5..286b0bf 100644 --- a/src/test/resources/application.conf +++ b/src/test/resources/application.conf @@ -1,12 +1,12 @@ # -# check the reference.conf in .../akka-actor_2.12-2.5.7.jar!/reference.conf for all defined settings +# check the reference.conf in .../pekko-actor_2.13-1.0.x.jar!/reference.conf for all defined settings # -akka { +pekko { # - # Akka version, checked against the runtime version of Akka. + # Pekko version, checked against the runtime version of Pekko. # - version = "2.5.29" + version = "1.0.2" loglevel = "DEBUG" diff --git a/src/test/scala/ai/agnos/sparql/SparqlQueries.scala b/src/test/scala/ai/agnos/sparql/SparqlQueries.scala index 2ffdf9b..23c4c73 100644 --- a/src/test/scala/ai/agnos/sparql/SparqlQueries.scala +++ b/src/test/scala/ai/agnos/sparql/SparqlQueries.scala @@ -3,12 +3,12 @@ package ai.agnos.sparql import ai.agnos.sparql.api._ import org.eclipse.rdf4j.model.IRI -import akka.http.scaladsl.model.HttpMethods._ +import org.apache.pekko.http.scaladsl.model.HttpMethods._ import ai.agnos.sparql.stream.client.SparqlClientConstants.{valueFactory => svf} trait SparqlQueries { - implicit val pm = PrefixMapping.extended + implicit val pm: PrefixMapping = PrefixMapping.extended implicit def stringToIri(iri: String): IRI = { svf.createIRI(iri.toString) diff --git a/src/test/scala/ai/agnos/sparql/stream/GraphStoreProtocolBuilderSpec.scala b/src/test/scala/ai/agnos/sparql/stream/GraphStoreProtocolBuilderSpec.scala index 640f4c3..84629f1 100644 --- a/src/test/scala/ai/agnos/sparql/stream/GraphStoreProtocolBuilderSpec.scala +++ b/src/test/scala/ai/agnos/sparql/stream/GraphStoreProtocolBuilderSpec.scala @@ -3,16 +3,16 @@ package ai.agnos.sparql.stream import java.io.File import java.net.URL -import akka.actor.ActorSystem -import akka.http.scaladsl.Http -import akka.http.scaladsl.Http.ServerBinding -import akka.http.scaladsl.model.{HttpEntity, HttpResponse} -import akka.http.scaladsl.server.Directives._ -import akka.stream.ActorMaterializer -import akka.stream.scaladsl._ -import akka.stream.testkit.TestSubscriber.Probe -import akka.stream.testkit.scaladsl.{TestSink, TestSource} -import akka.testkit.TestKit +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.Http.ServerBinding +import org.apache.pekko.http.scaladsl.model.{HttpEntity, HttpResponse} +import org.apache.pekko.http.scaladsl.server.Directives._ +import org.apache.pekko.stream.Materializer +import org.apache.pekko.stream.scaladsl._ +import org.apache.pekko.stream.testkit.TestSubscriber.Probe +import org.apache.pekko.stream.testkit.scaladsl.{TestSink, TestSource} +import org.apache.pekko.testkit.TestKit import scala.concurrent.{Await, ExecutionContext, Future} import scala.concurrent.duration._ @@ -39,11 +39,11 @@ class GraphStoreProtocolBuilderSpec extends TestKit(ActorSystem("GraphStoreProto with SparqlQueries with RdfModelTestUtils { - implicit val materializer: ActorMaterializer = ActorMaterializer()(system) + implicit val materializer: Materializer = Materializer(system) implicit val dispatcher: ExecutionContext = system.dispatcher implicit val prefixMapping: PrefixMapping = PrefixMapping.none - import scala.collection.JavaConverters._ + import scala.jdk.CollectionConverters._ implicit val errorHandler: ErrorHandler = DefaultErrorHandler @@ -128,6 +128,8 @@ class GraphStoreProtocolBuilderSpec extends TestKit(ActorSystem("GraphStoreProto } case SparqlResponse(request, success, _, result, error) => info(s"response status for $request ===>>> $success / ${result.size} items/ error: $error") + case _ => + fail(s"unexpected response: $response") } response @@ -174,6 +176,7 @@ class GraphStoreProtocolBuilderSpec extends TestKit(ActorSystem("GraphStoreProto checkAllGood(sparqlSink) match { case SparqlResponse (_, true, _, result, None) => assert(result === query1Result) + case _ => fail("unexpected response") } // now check with GetGraph (graph-store protocol variant) @@ -194,6 +197,7 @@ class GraphStoreProtocolBuilderSpec extends TestKit(ActorSystem("GraphStoreProto checkAllGood(sparqlSink) match { case SparqlResponse (_, true, _, result, None) => assert(result === query2Result) + case _ => fail("unexpected response") } // now check with GetGraph (graph-store protocol variant) @@ -215,6 +219,7 @@ class GraphStoreProtocolBuilderSpec extends TestKit(ActorSystem("GraphStoreProto checkAllGood(sparqlSink) match { case SparqlResponse (_, true, _, result, None) => assert(result === query2Result) + case _ => fail("unexpected response") } // now check with GetGraph (graph-store protocol variant) @@ -299,10 +304,12 @@ class GraphStoreProtocolBuilderSpec extends TestKit(ActorSystem("GraphStoreProto case path if path.toString == "labels.ttl" => val fileSource = FileIO.fromPath(new File(s"$rootFolder/$path").toPath) complete(HttpResponse(entity = HttpEntity(`text/turtle`, fileSource))) + case path => + complete(HttpResponse(status = 404, entity = s"file not found: $path") ) } } - Http().bindAndHandle(route, serverEndpoint.host, serverEndpoint.port, log = system.log) + Http().newServerAt(serverEndpoint.host, serverEndpoint.port).bind(route) } val endpoint = HttpEndpoint.localhostWithAutomaticPort("/resources") diff --git a/src/test/scala/ai/agnos/sparql/stream/MappingStreamSparqlClientSpec.scala b/src/test/scala/ai/agnos/sparql/stream/MappingStreamSparqlClientSpec.scala index 1acd6fa..869b706 100644 --- a/src/test/scala/ai/agnos/sparql/stream/MappingStreamSparqlClientSpec.scala +++ b/src/test/scala/ai/agnos/sparql/stream/MappingStreamSparqlClientSpec.scala @@ -1,20 +1,19 @@ package ai.agnos.sparql.stream -import akka.actor.ActorSystem -import akka.http.scaladsl.Http -import akka.stream.ActorMaterializer -import akka.stream.scaladsl.Keep -import akka.stream.testkit.scaladsl.{TestSink, TestSource} -import akka.testkit.TestKit +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.stream.Materializer +import org.apache.pekko.stream.scaladsl.Keep +import org.apache.pekko.stream.testkit.scaladsl.{TestSink, TestSource} +import org.apache.pekko.testkit.TestKit import ai.agnos.sparql.SparqlQueries import ai.agnos.sparql.api._ import ai.agnos.sparql.stream.client.{HttpClientFlowBuilder, HttpEndpointFlow, SparqlRequestFlowBuilder} import ai.agnos.test.HttpEndpointSuiteTestRunner import org.scalatest._ -import scala.concurrent.Await +import scala.concurrent.{Await, ExecutionContext} import scala.concurrent.duration._ -import scala.language.postfixOps /** * This test runs as part of the [[HttpEndpointSuiteTestRunner]] Suite. @@ -23,13 +22,13 @@ import scala.language.postfixOps class MappingStreamSparqlClientSpec() extends TestKit(ActorSystem("MappingStreamSparqlClientSpec")) with WordSpecLike with MustMatchers with BeforeAndAfterAll with SparqlQueries with SparqlRequestFlowBuilder with HttpClientFlowBuilder { - implicit val materializer = ActorMaterializer()(system) - implicit val dispatcher = system.dispatcher - implicit val prefixMapping = PrefixMapping.none + implicit val materializer: Materializer = Materializer(system) + implicit val dispatcher: ExecutionContext = system.dispatcher + implicit val prefixMapping: PrefixMapping = PrefixMapping.none implicit val errorHandler: ErrorHandler = DefaultErrorHandler - val receiveTimeout = 5 seconds + val receiveTimeout: FiniteDuration = 5 seconds import HttpEndpointSuiteTestRunner._ @@ -52,6 +51,7 @@ class MappingStreamSparqlClientSpec() extends TestKit(ActorSystem("MappingStream sink.expectNext(receiveTimeout) match { case SparqlResponse (_, true, _, result, None) => assert(result === emptyResult) + case _ => fail("unexpected response") } } diff --git a/src/test/scala/ai/agnos/sparql/stream/SparqlConstructClientSpec.scala b/src/test/scala/ai/agnos/sparql/stream/SparqlConstructClientSpec.scala index a251e00..12a3c8c 100644 --- a/src/test/scala/ai/agnos/sparql/stream/SparqlConstructClientSpec.scala +++ b/src/test/scala/ai/agnos/sparql/stream/SparqlConstructClientSpec.scala @@ -1,12 +1,12 @@ package ai.agnos.sparql.stream -import akka.actor.ActorSystem -import akka.http.scaladsl.Http -import akka.stream.ActorMaterializer -import akka.stream.scaladsl.Keep -import akka.stream.testkit.{TestPublisher, TestSubscriber} -import akka.stream.testkit.scaladsl.{TestSink, TestSource} -import akka.testkit.TestKit +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.stream.Materializer +import org.apache.pekko.stream.scaladsl.Keep +import org.apache.pekko.stream.testkit.{TestPublisher, TestSubscriber} +import org.apache.pekko.stream.testkit.scaladsl.{TestSink, TestSource} +import org.apache.pekko.testkit.TestKit import ai.agnos.sparql.SparqlQueries import ai.agnos.sparql.api._ import ai.agnos.sparql.stream.client.{HttpClientFlowBuilder, HttpEndpointFlow, SparqlRequestFlowBuilder} @@ -16,7 +16,6 @@ import org.scalatest._ import scala.concurrent.duration._ import scala.concurrent.{Await, ExecutionContext} -import scala.language.postfixOps import scala.util.{Failure, Success, Try} /** @@ -27,7 +26,7 @@ class SparqlConstructClientSpec extends TestKit(ActorSystem("SparqlToModelConstructClientSpec")) with SparqlConstructSpecBase { - implicit val materializer: ActorMaterializer = ActorMaterializer()(system) + implicit val materializer: Materializer = Materializer(system) implicit val dispatcher: ExecutionContext = system.dispatcher implicit val prefixMapping: PrefixMapping = PrefixMapping.none @@ -234,6 +233,7 @@ trait SparqlConstructSpecBase source.sendNext(SparqlRequest(queryModelGraph)) sink.expectNext(receiveTimeout) match { case SparqlResponse (_, true, _, result, None) => assert(result == emptyResult) + case _ => fail("unexpected response") } } diff --git a/src/test/scala/ai/agnos/sparql/stream/SparqlHttpSpec.scala b/src/test/scala/ai/agnos/sparql/stream/SparqlHttpSpec.scala index 5ae27c9..77f14e2 100644 --- a/src/test/scala/ai/agnos/sparql/stream/SparqlHttpSpec.scala +++ b/src/test/scala/ai/agnos/sparql/stream/SparqlHttpSpec.scala @@ -1,23 +1,22 @@ package ai.agnos.sparql.stream -import akka.actor.ActorSystem +import org.apache.pekko.actor.ActorSystem -import akka.stream.ActorMaterializer -import akka.stream.scaladsl.Flow -import akka.stream.scaladsl._ +import org.apache.pekko.stream.Materializer +import org.apache.pekko.stream.scaladsl.Flow +import org.apache.pekko.stream.scaladsl._ -import akka.http.scaladsl.model._ -import akka.http.scaladsl.Http -import akka.http.scaladsl.model.headers.{Authorization, BasicHttpCredentials} +import org.apache.pekko.http.scaladsl.model._ +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.model.headers.{Authorization, BasicHttpCredentials} import ai.agnos.test.HttpEndpointSuiteTestRunner import scala.concurrent.{Await, Future} import scala.concurrent.duration._ -import scala.language.postfixOps import org.scalatest._ -import akka.testkit.TestKit +import org.apache.pekko.testkit.TestKit /** * This test runs as part of the [[HttpEndpointSuiteTestRunner]] Suite. @@ -26,7 +25,7 @@ import akka.testkit.TestKit class SparqlHttpSpec extends TestKit(ActorSystem("SparqlHttpSpec")) with WordSpecLike with MustMatchers with BeforeAndAfterAll { - implicit val testMaterializer = ActorMaterializer() + implicit val testMaterializer: Materializer = Materializer(system) import HttpEndpointSuiteTestRunner.testServerEndpoint._ diff --git a/src/test/scala/ai/agnos/sparql/stream/SparqlQueryConstructionSpec.scala b/src/test/scala/ai/agnos/sparql/stream/SparqlQueryConstructionSpec.scala index 47838f8..a93eaad 100644 --- a/src/test/scala/ai/agnos/sparql/stream/SparqlQueryConstructionSpec.scala +++ b/src/test/scala/ai/agnos/sparql/stream/SparqlQueryConstructionSpec.scala @@ -2,7 +2,7 @@ package ai.agnos.sparql.stream import ai.agnos.sparql.api.SparqlQuery import org.scalatest.WordSpec -import akka.http.scaladsl.model.HttpMethods._ +import org.apache.pekko.http.scaladsl.model.HttpMethods._ class SparqlQueryConstructionSpec extends WordSpec { diff --git a/src/test/scala/ai/agnos/sparql/stream/SparqlRequestClientSpec.scala b/src/test/scala/ai/agnos/sparql/stream/SparqlRequestClientSpec.scala index 028a8ab..0b5f9ce 100644 --- a/src/test/scala/ai/agnos/sparql/stream/SparqlRequestClientSpec.scala +++ b/src/test/scala/ai/agnos/sparql/stream/SparqlRequestClientSpec.scala @@ -1,11 +1,11 @@ package ai.agnos.sparql.stream -import akka.actor.ActorSystem -import akka.http.scaladsl.Http -import akka.stream.ActorMaterializer -import akka.stream.scaladsl.Keep -import akka.stream.testkit.scaladsl.{TestSink, TestSource} -import akka.testkit.TestKit +import org.apache.pekko.actor.ActorSystem +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.stream.Materializer +import org.apache.pekko.stream.scaladsl.Keep +import org.apache.pekko.stream.testkit.scaladsl.{TestSink, TestSource} +import org.apache.pekko.testkit.TestKit import ai.agnos.sparql.SparqlQueries import ai.agnos.sparql.api._ import ai.agnos.sparql.stream.client.{HttpClientFlowBuilder, HttpEndpointFlow, SparqlRequestFlowBuilder} @@ -14,7 +14,6 @@ import org.scalatest._ import scala.concurrent.{Await, ExecutionContext} import scala.concurrent.duration._ -import scala.language.postfixOps /** * This test runs as part of the [[HttpEndpointSuiteTestRunner]] Suite. @@ -24,7 +23,7 @@ class SparqlRequestClientSpec extends TestKit(ActorSystem("SparqlRequestClientSp with WordSpecLike with MustMatchers with BeforeAndAfterAll with SparqlQueries with SparqlRequestFlowBuilder with HttpClientFlowBuilder { - implicit val materializer: ActorMaterializer = ActorMaterializer()(system) + implicit val materializer: Materializer = Materializer(system) implicit val dispatcher: ExecutionContext = system.dispatcher implicit val prefixMapping: PrefixMapping = PrefixMapping.none @@ -91,6 +90,7 @@ class SparqlRequestClientSpec extends TestKit(ActorSystem("SparqlRequestClientSp sink.expectNext(receiveTimeout) match { case SparqlResponse (_, true, _, result, None) => assert(result === query2Result) + case _ => fail("unexpected response") } } @@ -101,6 +101,7 @@ class SparqlRequestClientSpec extends TestKit(ActorSystem("SparqlRequestClientSp sink.expectNext(receiveTimeout) match { case SparqlResponse (_, true, _, result, None) => assert(result === query2Result) + case _ => fail("unexpected response") } } diff --git a/src/test/scala/ai/agnos/test/FusekiManager.scala b/src/test/scala/ai/agnos/test/FusekiManager.scala index f467b40..a296339 100644 --- a/src/test/scala/ai/agnos/test/FusekiManager.scala +++ b/src/test/scala/ai/agnos/test/FusekiManager.scala @@ -1,10 +1,10 @@ package ai.agnos.test -import akka.actor._ -import akka.http.scaladsl.Http -import akka.http.scaladsl.model.{HttpMethods, HttpRequest, HttpResponse, StatusCodes} -import akka.stream.ActorMaterializer -import akka.util.Timeout +import org.apache.pekko.actor._ +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.http.scaladsl.model.{HttpMethods, HttpRequest, HttpResponse, StatusCodes} +import org.apache.pekko.stream.Materializer +import org.apache.pekko.util.Timeout import ai.agnos.sparql.util.HttpEndpoint import ai.agnos.test.FusekiManager._ @@ -75,9 +75,9 @@ object FusekiManager { class FusekiManager(val endpoint: HttpEndpoint) extends Actor with ActorLogging { import context.dispatcher - implicit val system = context.system - implicit val materializer = ActorMaterializer() - implicit val timeout = Timeout(5 seconds) + implicit val system: ActorSystem = context.system + implicit val materializer: Materializer = Materializer(context.system) + implicit val timeout: Timeout = Timeout(5 seconds) private val fusekiRunner = new FusekiRunner(endpoint.port, endpoint.path) @@ -102,7 +102,7 @@ class FusekiManager(val endpoint: HttpEndpoint) extends Actor with ActorLogging // with multiple parameters to control behavior, not straightforward to understand //SSZ: True, but again, this is a test class that does the work for us already. I would not worry // about making it more understandable, unless you really think it would add more value? - self ! Ping(sender, sendOnSuccess = StartOk, sendOnFailure = StartError, pingInterval = 3 seconds, retriesLeft = 30) + self ! Ping(sender(), sendOnSuccess = StartOk, sendOnFailure = StartError, pingInterval = 3 seconds, retriesLeft = 30) case x @ Ping(originalSender, sendOnSuccess, sendOnFailure, pingInterval, stopOnSuccess, 0) => log.info(s"Ping timeout for $x") @@ -133,7 +133,7 @@ class FusekiManager(val endpoint: HttpEndpoint) extends Actor with ActorLogging } case Shutdown => - val originalSender = sender + val originalSender = sender() log.info(s"Sending: $shutdownReq") pipeline(shutdownReq) onComplete { case Success(HttpResponse(StatusCodes.OK, _, _, _)) => diff --git a/src/test/scala/ai/agnos/test/HttpEndpointSuiteTestRunner.scala b/src/test/scala/ai/agnos/test/HttpEndpointSuiteTestRunner.scala index b612e2e..2fadbde 100644 --- a/src/test/scala/ai/agnos/test/HttpEndpointSuiteTestRunner.scala +++ b/src/test/scala/ai/agnos/test/HttpEndpointSuiteTestRunner.scala @@ -1,8 +1,8 @@ package ai.agnos.test -import akka.actor.{ActorSystem, Props} -import akka.http.scaladsl.Http -import akka.testkit.{ImplicitSender, TestKit} +import org.apache.pekko.actor.{ActorSystem, Props} +import org.apache.pekko.http.scaladsl.Http +import org.apache.pekko.testkit.{ImplicitSender, TestKit} import ai.agnos.sparql.stream._ import ai.agnos.sparql.util.{BasicAuthentication, HttpEndpoint} import ai.agnos.test.FusekiManager._ @@ -35,15 +35,15 @@ object HttpEndpointSuiteTestRunner { val config: Config = { ConfigFactory.parseString( s""" - |akka.loggers = ["akka.testkit.TestEventListener"] - |akka.loglevel = INFO - |akka.remote { + |pekko.loggers = ["org.apache.pekko.testkit.TestEventListener"] + |pekko.loglevel = INFO + |pekko.remote { | netty.tcp { | hostname = "" | port = 0 | } |} - |akka.cluster { + |pekko.cluster { | seed-nodes = [] |} |sparql.client { @@ -52,7 +52,7 @@ object HttpEndpointSuiteTestRunner { | password = "admin" |} | - |akka { + |pekko { | http { | server.parsing.illegal-header-warnings = off | client.parsing.illegal-header-warnings = off @@ -98,21 +98,21 @@ class HttpEndpointSuiteTestRunner(_system: ActorSystem) extends TestKit(_system) new GraphStoreProtocolBuilderSpec() ) - val _log = akka.event.Logging(this.system, testActor) + val _log = org.apache.pekko.event.Logging(this.system, testActor) private lazy val fusekiManager = system.actorOf(Props(classOf[FusekiManager], testServerEndpoint), "fuseki-manager") val startTimeout = 100 seconds val stopTimeout = 20 seconds - override def beforeAll() { + override def beforeAll(): Unit = { if (useFuseki) { fusekiManager ! Start expectMsg(startTimeout, StartOk) } } - override def afterAll() { + override def afterAll(): Unit = { if (useFuseki) { fusekiManager ! Shutdown expectMsg(stopTimeout, "Allowing Fuseki Server to shut down", ShutdownOk) @@ -125,7 +125,7 @@ class HttpEndpointSuiteTestRunner(_system: ActorSystem) extends TestKit(_system) * * @param system The ActorSystem currently in use. */ - def shutdownSystem(implicit system: ActorSystem) { + def shutdownSystem(implicit system: ActorSystem): Unit = { Await.result(Http().shutdownAllConnectionPools(), 5 seconds) Await.result(system.terminate(), 5 seconds) }