Ir para o conteúdo

Concorrência ​

O Talor não tem async nem await. A concorrência são duas bibliotecas: std.fiber roda muitas tasks leves em poucas threads, e std.thread inicia threads do sistema operacional. Uma função que lê um socket tem uma assinatura só: dentro de uma fiber a leitura estaciona a fiber e a thread roda outra, e fora de uma fiber ela bloqueia a thread. Nenhuma função ganha uma cor pelo jeito como é chamada.

O que a linguagem contribui são as regras de ownership que tornam isso seguro: uma task recebe os seus valores por move, o que duas tasks compartilham é read-only, e uma escrita de que duas tasks precisam passa por um lock que é dono dos dados.

Fibers ​

Uma fiber é uma stack própria e um contexto salvo. run(body) inicia o scheduler na thread que chamou e retorna quando todas as fibers que ele iniciou terminaram. Dentro dele, spawn(body) inicia uma fiber e join(f) espera por uma.

talor
use std.fiber.{run, spawn, join, sleep_ms};

fn main(): i32 {
    let done = run(() => {
        let slow = spawn(() => {
            match sleep_ms(20) { Ok(_) => {}, Err(_) => {} }
            println("the slow fiber ends second");
        });
        let fast = spawn(() => { println("the fast fiber ends first"); });
        match slow { Ok(f) => { match join(f) { Ok(_) => {}, Err(_) => {} } }, Err(_) => {} }
        match fast { Ok(f) => { match join(f) { Ok(_) => {}, Err(_) => {} } }, Err(_) => {} }
    });
    match done { Ok(_) => 0, Err(_) => 1 }
}

O resto de std.fiber:

ChamadaO que faz
run_on(n, body)roda as fibers em n threads, a que chamou e n - 1 que o runtime inicia; uma fiber pronta pode ser pega por uma thread que não tem nada para rodar
sleep_ms(ms), yield_now()estacionam a fiber por um tempo, ou deixam as outras rodarem
timeout(ms, body)roda body e responde Err (código timed_out_code()) se ele não terminou em ms
cancel(f), cancelled()pedem que uma fiber pare: ela ouve o pedido onde esperaria (sleep_ms, yield_now, join), que responde Err ali, ou lendo cancelled(); nada a para de fora, e as fibers que ela iniciou continuam rodando
offload(body)roda body em uma thread de um pool e estaciona a fiber até ele retornar: é o jeito de chamar código C que bloqueia sem parar as outras fibers
all(bodies)roda uma lista de closures ao mesmo tempo e responde os valores delas na ordem da lista (abaixo)

A stack de uma fiber não se move, então uma view que ela segura continua válida depois de um estacionamento. Uma fiber que passa do fim da sua stack encerra o programa com uma mensagem que nomeia a fiber e o tamanho da stack, em vez de escrever na memória de outra fiber.

all: fibers que pegam por borrow o que capturam ​

Uma closure escrita diretamente na lista passada a fiber.all pega o que captura por borrow durante a chamada, em vez de por move, e quem chamou recebe os valores de volta, já escritos, quando all retorna. Duas closures da lista não podem tocar o mesmo lugar quando uma delas o escreve; escrever dois campos diferentes é aceito:

talor
use std.fiber;

struct Halves { left: Array<i64>, right: Array<i64> }

fn main(): i32 {
    let done = fiber.run(() => {
        let t = Halves { left: [1, 2, 3], right: [4, 5, 6] };
        let lens = fiber.all([
            () => { t.left.push(0); t.left.len() },
            () => { t.right.push(0); t.right.len() },
        ]);
        match lens {
            Ok(ls) => println(`lengths ${ls[0]} and ${ls[1]}`),
            Err(e) => println(`failed: ${e.message()}`),
        }
        println(`left ends in ${t.left[3]}, right ends in ${t.right[3]}`);
    });
    match done { Ok(_) => 0, Err(_) => 1 }
}

t.left em uma closure e t.right na outra são disjuntos; t em uma e t.left na outra seriam recusados, e xs[0] e xs[1] também, porque um índice não é diferenciado de outro. Uma closure da lista não pode repassar uma captura (devolvê-la, guardá-la, ou passá-la a um parâmetro que fica com ela), já que a captura está só em borrow. Quando uma delas devolve um erro ou entra em panic, all responde a primeira falha na ordem da lista, depois que todas as closures terminaram.

Threads ​

std.thread inicia uma closure em uma thread própria com start, e join(t) espera por ela e detach(t) a deixa ir. hardware_count() diz quantas threads a máquina roda ao mesmo tempo.

A closure passada a spawn ou start pega as suas capturas por move, então a task é dona do que recebeu. Ela pode ser uma closure que é chamada uma vez, porque cada um roda a sua closure exatamente uma vez. Uma closure que a função chamada guarda não pode capturar uma view, já que a task pode viver mais que o frame de onde a view veio.

Compartilhar: shared T ​

Um valor de que duas tasks precisam é colocado em um shared, e cada task recebe um clone dele: um clone é uma contagem a mais na mesma alocação, e não uma cópia do valor. Um lugar alcançado através de um shared é read-only, porque duas tasks poderiam escrevê-lo no mesmo instante; uma escrita através de um deles é recusada em tempo de compilação. A contagem é simples até o valor poder ser alcançado por uma segunda thread, e atomic dali em diante, então um programa que não compartilha nada não paga nada por isso.

Channels ​

Um Channel<T> leva valores de qualquer tipo de uma task a outra, inteiros: send fica com o valor e recv o entrega a quem recebe.

talor
use std.thread.{start, join, channel, send, recv, close};

fn main(): i32 {
    let ch = match channel<string>(4) { Ok(c) => c, Err(_) => { return 1; } };
    let owner = shared ch;
    let mine = owner.clone();
    let worker = start(() => {
        for word in ["one", "two", "three"] {
            match send(*mine, word) { Ok(_) => {}, Err(_) => {} }
        }
        match close(*mine) { Ok(_) => {}, Err(_) => {} }
    });
    let going = true;
    while going {
        match recv(*owner) {
            Ok(Some(word)) => println(`received ${word}`),
            Ok(None) => { going = false; }
            Err(_) => { going = false; }
        }
    }
    match worker { Ok(t) => { match join(t) { Ok(_) => 0, Err(_) => 1 } }, Err(_) => 1 }
}
  • channel<T>(capacity) cria um. Uma capacidade positiva faz send esperar enquanto o channel está cheio; 0 quer dizer sem limite.
  • recv espera por um valor e responde None quando o channel está fechado e vazio; try_recv nunca espera.
  • send em um channel fechado responde Err, e o valor que ele recebeu é liberado ali.
  • select(chans, timeout_ms) espera em vários channels e responde qual deles tinha um valor.
  • Um channel liberado com valores ainda dentro libera cada um deles.

Dentro de uma fiber, esperar em um channel estaciona a fiber; em uma thread, bloqueia a thread.

Locks são donos do que protegem ​

Um lock não é um valor parado ao lado dos dados: mutex(v) fica com v, e o único caminho até os dados é with(m, body), que pega o lock, empresta os dados a body e libera o lock quando body retorna. O parâmetro do body é um Locked<T>, um lugar que pode ser escrito durante essa única chamada e não pode ser guardado, devolvido nem capturado.

talor
use std.thread.{start, join, mutex, with, Thread};

struct Counter { n: i64 }

fn main(): i32 {
    let m = shared mutex(Counter { n: 0 });
    let threads: Array<Thread> = [];
    for i in 0..4 {
        let mine = m.clone();
        match start(() => {
            for k in 0..1000 {
                match with(*mine, (c) => { c.n += 1; }) { Ok(_) => {}, Err(_) => {} }
            }
        }) {
            Ok(t) => threads.push(t),
            Err(_) => { return 1; }
        }
    }
    while threads.len() > 0 {
        match threads.pop() { Some(t) => { match join(t) { Ok(_) => {}, Err(_) => {} } }, None => {} }
    }
    match with(*m, (c) => c.n) {
        Ok(n) => { println(`four threads counted ${n}`); 0 }
        Err(_) => 1,
    }
}
  • try_with(m, body) roda body só se o lock estiver livre, e responde None sem esperar caso contrário.
  • rwlock(v) admite muitos leitores ou um escritor: read(l, body) empresta os dados como read-only a qualquer número de leitores ao mesmo tempo, e write(l, body) empresta um Locked<T> a um só.
  • Uma função auxiliar pode receber o lugar do borrow pelo seu tipo plain, fn bump(c: Counter), e ser chamada de dentro do body.
  • Uma task que pega um lock que já segura é parada com um panic em vez de ficar esperando para sempre, e um lock cujo dono entrou em panic fica envenenado: o próximo with responde Err.

Atomics ​

Para um único contador ou flag, um Atomic sai mais barato que um lock: atomic(n) cria um, e load, store, add e compare_exchange o leem e escrevem em um passo só.

talor
use std.thread.{start, join, atomic, add, load, Thread};

fn main(): i32 {
    let hits = match atomic(0) { Ok(a) => shared a, Err(_) => { return 1; } };
    let threads: Array<Thread> = [];
    for i in 0..4 {
        let mine = hits.clone();
        match start(() => { for k in 0..1000 { add(*mine, 1); } }) {
            Ok(t) => threads.push(t),
            Err(_) => { return 1; }
        }
    }
    while threads.len() > 0 {
        match threads.pop() { Some(t) => { match join(t) { Ok(_) => {}, Err(_) => {} } }, None => {} }
    }
    println(`hits: ${load(*hits)}`);
    0
}

Quando uma task entra em panic ​

Um panic encerra a fiber ou a thread em que acontece, e não o processo: os frames da task liberam o que possuem no caminho de saída, PANIC: e o local são escritos na saída de erro padrão, e quem faz join da task recebe Err com código 101, o texto do panic, e e.panicked() verdadeiro. Um panic na main thread, ou no body passado a run, encerra o processo com código 101.

talor
use std.fiber.{run, spawn, join};

fn main(): i32 {
    let done = run(() => {
        match spawn(() => { let xs: Array<i64> = []; println(`${xs[3]}`); }) {
            Ok(f) => {
                match join(f) {
                    Ok(_) => println("joined"),
                    Err(e) => println(`the fiber panicked: ${e.panicked()}, code ${e.code()}`),
                }
            }
            Err(_) => {}
        }
        println("the run goes on");
    });
    match done { Ok(_) => 0, Err(_) => 1 }
}

Talor v0.1.0 - Publicado sob a licença MIT OR Apache-2.0.