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.
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:
| Chamada | O 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:
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.
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 fazsendesperar enquanto o channel está cheio;0quer dizer sem limite.recvespera por um valor e respondeNonequando o channel está fechado e vazio;try_recvnunca espera.sendem um channel fechado respondeErr, 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.
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)rodabodysó se o lock estiver livre, e respondeNonesem 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, ewrite(l, body)empresta umLocked<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
withrespondeErr.
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ó.
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.
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 }
}