-
Notifications
You must be signed in to change notification settings - Fork 641
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Directory sources #274
Directory sources #274
Changes from 1 commit
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,26 @@ | ||
/* | ||
* Copyright (C) 2016-2017 Lightbend Inc. <http://www.lightbend.com> | ||
*/ | ||
package akka.stream.alpakka.file.javadsl; | ||
|
||
import akka.NotUsed; | ||
import akka.stream.javadsl.Source; | ||
|
||
import java.nio.file.Path; | ||
|
||
public final class Directory { | ||
|
||
/** | ||
* List all files in the given directory | ||
*/ | ||
public static Source<Path, NotUsed> ls(Path directory) { | ||
return akka.stream.alpakka.file.scaladsl.Directory.ls(directory).asJava(); | ||
} | ||
|
||
/** | ||
* Recursively list files and directories in the given directory, depth first. | ||
*/ | ||
public static Source<Path, NotUsed> walk(Path directory) { | ||
return akka.stream.alpakka.file.scaladsl.Directory.walk(directory).asJava(); | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,36 @@ | ||
/* | ||
* Copyright (C) 2016-2017 Lightbend Inc. <http://www.lightbend.com> | ||
*/ | ||
package akka.stream.alpakka.file.scaladsl | ||
|
||
import java.nio.file.attribute.BasicFileAttributes | ||
import java.nio.file.{Files, Path} | ||
import java.util.function.BiPredicate | ||
|
||
import akka.NotUsed | ||
import akka.stream.scaladsl.{Source, StreamConverters} | ||
|
||
object Directory { | ||
|
||
/** | ||
* List all files in the given directory | ||
*/ | ||
def ls(directory: Path): Source[Path, NotUsed] = { | ||
require(Files.isDirectory(directory), s"Path must be a directory, $directory isn't") | ||
StreamConverters.fromJavaStream(() => Files.list(directory)) | ||
} | ||
|
||
/** | ||
* Recursively list files in the given directory and its subdirectories. Listing is done | ||
* depth first. | ||
*/ | ||
def walk(directory: Path): Source[Path, NotUsed] = { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think it would be good to expose the overloads too: There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I started out with that, but then thought that it is so simple to create it yourself, so maybe more specific needs can be met by just looking at the sources of this and copying that single line of code... There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. I think with alpakka we're targeting an "works out of the box" and "I would not look at the sources" use cases/people, so I'd really add those parameters IMHO There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This was rather from the not covering every conceivable use case under the sun, just the "normal" ones point of view, especially the filter which could be achieved by composition (with a bit of alloc overhead, but hey this is disk io we are talking about here). Could you come up with a common use case for depths other than 0, which is covered by ls and all the depth ( There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Depth is pretty important to protect from entering infinite loops with symlinks to an upper directory if I remember correctly how these APIs work? There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Alright, I'll add it. Skipping filter though, as that can be achieved with composition. |
||
require(Files.isDirectory(directory), s"Path must be a directory, $directory isn't") | ||
StreamConverters.fromJavaStream(() => Files.walk(directory)) | ||
} | ||
|
||
private def allFilesFilter = new BiPredicate[Path, BasicFileAttributes] { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. this is not used |
||
override def test(t: Path, u: BasicFileAttributes): Boolean = true | ||
} | ||
|
||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,92 @@ | ||
/* | ||
* Copyright (C) 2016-2017 Lightbend Inc. <http://www.lightbend.com> | ||
*/ | ||
package akka.stream.alpakka.file.javadsl; | ||
|
||
import akka.NotUsed; | ||
import akka.actor.ActorSystem; | ||
import akka.stream.ActorMaterializer; | ||
import akka.stream.Materializer; | ||
import akka.stream.javadsl.Sink; | ||
import akka.stream.javadsl.Source; | ||
import akka.testkit.TestKit; | ||
import com.google.common.jimfs.Configuration; | ||
import com.google.common.jimfs.Jimfs; | ||
import org.junit.After; | ||
import org.junit.Before; | ||
import org.junit.Test; | ||
import org.scalatest.junit.JUnitSuite; | ||
import scala.concurrent.duration.FiniteDuration; | ||
|
||
import java.nio.file.FileSystem; | ||
import java.nio.file.Files; | ||
import java.nio.file.Path; | ||
import java.util.Arrays; | ||
import java.util.List; | ||
import java.util.concurrent.TimeUnit; | ||
|
||
import static org.junit.Assert.assertEquals; | ||
|
||
public class DirectoryTest { | ||
|
||
private FileSystem fs; | ||
private ActorSystem system; | ||
private Materializer materializer; | ||
|
||
@Before | ||
public void setup() { | ||
fs = Jimfs.newFileSystem(Configuration.unix()); | ||
system = ActorSystem.create(); | ||
materializer = ActorMaterializer.create(system); | ||
} | ||
|
||
|
||
@Test | ||
public void listFiles() throws Exception { | ||
final Path dir = fs.getPath("listfiles"); | ||
Files.createDirectories(dir); | ||
final Path file1 = Files.createFile(dir.resolve("file1")); | ||
final Path file2 = Files.createFile(dir.resolve("file2")); | ||
|
||
// #ls | ||
final Source<Path, NotUsed> source = Directory.ls(dir); | ||
// #ls | ||
|
||
final List<Path> result = source.runWith(Sink.seq(), materializer) | ||
.toCompletableFuture().get(3, TimeUnit.SECONDS); | ||
assertEquals(result.size(), 2); | ||
assertEquals(result.get(0), file1); | ||
assertEquals(result.get(1), file2); | ||
} | ||
|
||
@Test | ||
public void walkAFileTree() throws Exception { | ||
final Path root = fs.getPath("walk"); | ||
Files.createDirectories(root); | ||
final Path subdir1 = root.resolve("subdir1"); | ||
Files.createDirectories(subdir1); | ||
final Path file1 = subdir1.resolve("file1"); | ||
Files.createFile(file1); | ||
final Path subdir2 = root.resolve("subdir2"); | ||
Files.createDirectories(subdir2); | ||
final Path file2 = subdir2.resolve("file2"); | ||
Files.createFile(file2); | ||
|
||
// #walk | ||
final Source<Path, NotUsed> source = Directory.walk(root); | ||
// #walk | ||
|
||
final List<Path> result = source.runWith(Sink.seq(), materializer) | ||
.toCompletableFuture().get(3, TimeUnit.SECONDS); | ||
assertEquals(result, Arrays.asList(root, subdir1, file1, subdir2, file2)); | ||
} | ||
|
||
@After | ||
public void tearDown() throws Exception { | ||
fs.close(); | ||
fs = null; | ||
TestKit.shutdownActorSystem(system, FiniteDuration.create(10, TimeUnit.SECONDS), true); | ||
system = null; | ||
materializer = null; | ||
} | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,67 @@ | ||
/* | ||
* Copyright (C) 2016-2017 Lightbend Inc. <http://www.lightbend.com> | ||
*/ | ||
package akka.stream.alpakka.file.scaladsl | ||
|
||
import java.nio.file.{Files, Path} | ||
|
||
import akka.NotUsed | ||
import akka.actor.ActorSystem | ||
import akka.stream.ActorMaterializer | ||
import akka.stream.scaladsl.{Sink, Source} | ||
import akka.testkit.TestKit | ||
import com.google.common.jimfs.{Configuration, Jimfs} | ||
import org.scalatest.concurrent.ScalaFutures | ||
import org.scalatest.{BeforeAndAfterAll, Matchers, WordSpecLike} | ||
|
||
class DirectorySpec | ||
extends TestKit(ActorSystem("directoryspec")) | ||
with WordSpecLike | ||
with Matchers | ||
with BeforeAndAfterAll | ||
with ScalaFutures { | ||
|
||
private val fs = Jimfs.newFileSystem(Configuration.forCurrentPlatform.toBuilder.build) | ||
private implicit val mat = ActorMaterializer() | ||
|
||
"The directory source factory" should { | ||
"list files" in { | ||
val dir = fs.getPath("listfiles") | ||
Files.createDirectories(dir) | ||
val paths = (0 to 100).map { n => | ||
val name = s"file$n" | ||
Files.createFile(dir.resolve(name)) | ||
} | ||
|
||
// #ls | ||
val source: Source[Path, NotUsed] = Directory.ls(dir) | ||
// #ls | ||
|
||
val result = source.runWith(Sink.seq).futureValue | ||
result.toSet shouldEqual paths.toSet | ||
} | ||
|
||
"walk a file tree" in { | ||
val root = fs.getPath("walk") | ||
Files.createDirectories(root) | ||
val subdir1 = root.resolve("subdir1") | ||
Files.createDirectories(subdir1) | ||
val file1 = subdir1.resolve("file1") | ||
Files.createFile(file1) | ||
val subdir2 = root.resolve("subdir2") | ||
Files.createDirectories(subdir2) | ||
val file2 = subdir2.resolve("file2") | ||
Files.createFile(file2) | ||
|
||
// #walk | ||
val files: Source[Path, NotUsed] = Directory.walk(root) | ||
// #walk | ||
|
||
val result = files.runWith(Sink.seq).futureValue | ||
result shouldEqual List(root, subdir1, file1, subdir2, file2) | ||
} | ||
} | ||
|
||
override protected def afterAll(): Unit = | ||
fs.close() | ||
} |
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -13,7 +13,7 @@ import akka.stream.{ActorMaterializer, Materializer} | |
|
||
import scala.concurrent.duration._ | ||
|
||
object FileTailSourceSpec extends { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. hah funny |
||
object FileTailSourceSpec { | ||
|
||
// small sample of usage, tails the first argument file path | ||
def main(args: Array[String]): Unit = { | ||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
files and directories