///|
struct Dataset[T] {
exec_ : async (MoonlightContext) -> @moonchor.Located[T, Manager]
}
///|
async fn[T] Dataset::exec(
self : Dataset[T],
ctx : MoonlightContext,
) -> @moonchor.Located[T, Manager] {
(self.exec_)(ctx)
}
///|
fn[T] Dataset::from_array(arr : Array[T]) -> Dataset[Array[T]] {
{
exec_: async fn(ctx) {
ctx.chorctx.locally(ctx.config.manager, fn(_unwrapper) { arr })
},
}
}
///|
fn[T : @moonchor.Message, U : @moonchor.Message] Dataset::map(
self : Dataset[Array[T]],
f : (T) -> U,
) -> Dataset[Array[U]] {
{
exec_: async fn(ctx) {
let data = self.exec(ctx)
let len = ctx.chorctx.locally_broadcast(ctx.config.manager, fn(
unwrapper,
) {
unwrapper.unwrap(data).length()
})
let results = ctx.chorctx.locally(ctx.config.manager, fn(_unwrapper) {
[]
})
for i in 0.. ignore
}
results
},
}
}
///|
pub async fn[T] Dataset::collect(
self : Dataset[Array[T]],
ctx : MoonlightContext,
) -> ImmediateData[Array[T]] {
let data = self.exec(ctx)
ImmediateData::new(data)
}
///|
async test "Dataset test" {
async fn dataset_test(ctx : MoonlightContext) {
let data = Dataset::from_array([1, 2, 3, 4])
let result = data.map(fn(x) { x * 1 })
let result = result.collect(ctx)
result.map(xs => println(Repr(xs))) |> ignore
}
let cluster = local_cluster_config(4)
cluster.start(dataset_test)
}