Skip to content
This repository has been archived by the owner. It is now read-only.
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Binary file modified .gradle/7.1/executionHistory/executionHistory.bin
Binary file not shown.
Binary file modified .gradle/7.1/executionHistory/executionHistory.lock
Binary file not shown.
Binary file modified .gradle/7.1/fileHashes/fileHashes.bin
Binary file not shown.
Binary file modified .gradle/7.1/fileHashes/fileHashes.lock
Binary file not shown.
Binary file modified .gradle/buildOutputCleanup/buildOutputCleanup.lock
Binary file not shown.
Binary file modified .gradle/buildOutputCleanup/outputFiles.bin
Binary file not shown.
6 changes: 6 additions & 0 deletions .idea/vcs.xml

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

122 changes: 114 additions & 8 deletions src/main/kotlin/Main.kt
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
import io.reactivex.rxjava3.core.Completable
import io.reactivex.rxjava3.core.Flowable
import io.reactivex.rxjava3.core.Maybe
import io.reactivex.rxjava3.core.Single
import java.rmi.server.ServerNotActiveException
import java.util.*
import java.util.concurrent.Callable
import java.util.concurrent.TimeUnit
import kotlin.random.Random

Expand All @@ -20,48 +22,152 @@ fun main() {
//
// Несмотря на то, что в некоторых заданиях фигурируют слова "синхронный" и "асинхронный" в рамках текущего ДЗ
// это всего лишь имитация, реальное переключение между потоками будет рассмотрено на следующем семинаре

println("Start Q1");

//Single.timer(1L,TimeUnit.MILLISECONDS).blockingSubscribe({println("0")},{println("error")});
//Some Emission
//Some Emission
//val singleSource = Maybe.just("single item");
/// maybe_array <ByteArray?>= Maybe.fromCallable ( {requestDataFromServerAsync();} );
/*
val count = -1
require(count >= 0) { println("Count must be non-negative, was $count") }
require(requestServerAsync() !is Unit) { println("kuku") }
//println("require(requestServerAsync() is Unit)"+require(requestServerAsync() is Unit));
*/
Maybe.fromCallable( {
// val arr: Array<ByteArray?> = arrayOfNulls<ByteArray?>(1);
//val arr: ByteArray? = arrayOfNulls<ByteArray?>(1)
// don't know standart fun like arrayOfNulls for List
// var result: MutableList<ByteArray?> = arr.toMutableList(); //arrayListOf<Int?>(list.size);
var result: ByteArray? = requestDataFromServerAsync()
// val result: <String> = "";
result
}).blockingSubscribe({s : ByteArray? -> println("Item received: from Maybe"+s.contentToString());
},
{ obj: Throwable -> obj.printStackTrace() } ) { println("Done from MaybeSource") }



/* singleSource.subscribe(
{ s: String -> println("Item received: from singleSource $s") },
{ obj: Throwable -> obj.printStackTrace() }
) { println("Done from SingleSource") }
*/


// Maybe<ByteArray?> maybe_array = Maybe.just("single item");
// Maybe<ByteArray?>.blockingSubscribe({},{});
// Maybe.fromCallable ( requestDataFromServerAsync() )
println("Finished Q1");


println("Start Q2");
//requestServerAsync();
//assertTrue( is Unit);
/* val callable1 = object : Callable<Int> {
override fun call(): Int = 1
}
val c1 = createCallable(callable1)
println("callable1 = ${c1.call()}")

val callable2 = object : Callable<Unit> {
override fun call(): Unit { println("Hello"); throw Exception("Hello") }
}
val c2 = createCallable(callable2)
c2.call()

*/

Completable.fromCallable ( object: Callable<Unit> {
override fun call(): Unit { requestServerAsync()}} ).
subscribe({println("Successful");},
{ obj: Throwable -> obj.printStackTrace() } );
//.fromCallable(requestServerAsync())
//.timer(1L,TimeUnit.MILLISECONDS).blockingSubscribe({println("2")},{println("error2")});


println("Finished Q2");

println("Start Q3");
//val val1 = requestDataFromDbAsync<Int?>();

/* Single.fromCallable ({
var result: Int? = requestDataFromDbAsync<Int?>() }
).onErrorComplete(println("eee"))
*/
//onErrorComplete(result: Throwable -> println("--"))

/* blockingSubscribe({println("Successful");},
{ obj: Throwable -> obj.printStackTrace() })

*/


/* ({s : Int? -> println("Item received: from Single:"+s.toString()})
{ obj: Throwable -> obj.printStackTrace() })
*/
println("Finish Q3");

println("Start Q4");
emitEachSecond();
println("Finish Q4");




}

// 1) Какой источник лучше всего подойдёт для запроса на сервер, который возвращает результат?
// Maybe
// Почему?
// Согласно лекции номер 2
// Дописать функцию
fun requestDataFromServerAsync() /* -> ???<ByteArray> */ {
fun requestDataFromServerAsync() : ByteArray? /* -> ???<ByteArray> */ {

// Функция имитирует синхронный запрос на сервер, возвращающий результат
fun getDataFromServerSync(): ByteArray? {
Thread.sleep(LATENCY);
val success = Random.nextBoolean()
return if (success) Random.nextBytes(RESPONSE_LENGTH) else null
}
//return if (success) Random.nextBytes(RESPONSE_LENGTH) else null
return Random.nextBytes(RESPONSE_LENGTH)

}
return getDataFromServerSync()
/* return ??? */
}


// 2) Какой источник лучше всего подойдёт для запроса на сервер, который НЕ возвращает результат?
//Completable
// Почему?
// По определению ..
// Дописать функцию
fun requestServerAsync() /* -> ??? */ {
fun requestServerAsync() :Unit /* -> ??? */ {

// Функция имитирует синхронный запрос на сервер, не возвращающий результат
fun getDataFromServerSync() {
Thread.sleep(LATENCY)
if (Random.nextBoolean()) throw ServerNotActiveException()
}

return getDataFromServerSync()
/* return ??? */
}

// 3) Какой источник лучше всего подойдёт для однократного асинхронного возвращения значения из базы данных?
//Single
// Почему?
//Single — реактивный Callable, потому что тут появляется возможность вернуть результат операции. Продолжая сравнение с Kotlin, можно сказать, что Single — это fun single(): T { }. Таким образом, чтобы подписаться на него, необходимо реализовать onSuccess(T) и onError.

// Дописать функцию
fun <T> requestDataFromDbAsync() /* -> ??? */ {
fun <T> requestDataFromDbAsync() : T? /* -> ??? */ {

// Функция имитирует синхронный запрос к БД не возвращающий результата
fun getDataFromDbSync(): T? {
Thread.sleep(LATENCY); return null
}

return getDataFromDbSync()
/* return */
}

Expand All @@ -76,7 +182,7 @@ fun emitEachSecond() {

// Принтер
fun printer(value: Long) = println("${Date()}: value = $value")

source().doOnEach { Thread.sleep(2000L)}.blockingSubscribe({printer(it)})
// code here
}

Expand Down