[refactor] move rbm.stack* + rbm.rnn_helper → stack/
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,32 @@
|
||||
from rbm.status import Status
|
||||
from rbm.train import train
|
||||
from stack.stack import Stack, StackType
|
||||
from rbm.matrix import Mat, np
|
||||
|
||||
class StackDeep(Stack):
|
||||
def __init__(self, name: str, work_dir: str = '.'):
|
||||
Stack.__init__(self, StackType.Deep, name, work_dir)
|
||||
|
||||
def train(self, batch: Mat, status=Status()):
|
||||
_batch = np.copy(batch)
|
||||
for index, layer in enumerate(self.layers):
|
||||
print(f"Train layer {index} for {layer.entity.training_params.num_epochs} epochs")
|
||||
train(layer.entity, _batch, status=status)
|
||||
_batch = layer.entity.forward(_batch)
|
||||
|
||||
def pass_up(self, visible: Mat, from_layer_id: int = 0):
|
||||
h = np.copy(visible)
|
||||
for layer in self.layers[from_layer_id:]:
|
||||
h = layer.entity.forward(h)
|
||||
return h
|
||||
|
||||
def pass_down(self, hidden: Mat, from_layer_id: int = 0):
|
||||
v = np.copy(hidden)
|
||||
for layer in list(reversed(self.layers))[from_layer_id:]:
|
||||
v = layer.entity.reconstruct(v)
|
||||
return v
|
||||
|
||||
def pass_down_up(self, visible: Mat, from_layer_id: int = 0):
|
||||
h = self.pass_up(visible, from_layer_id)
|
||||
v = self.pass_down(h)
|
||||
return v
|
||||
@@ -0,0 +1,71 @@
|
||||
import json
|
||||
from stack.stack import StackType
|
||||
from rbm.layer import Layer
|
||||
from rbm.entity import EntityParams, TrainingParams
|
||||
from stack.deep import StackDeep
|
||||
from stack.rnn import StackRnn
|
||||
|
||||
class StackFactory:
|
||||
@classmethod
|
||||
def from_dict(cls, project: dict, work_dir: str = ".") -> StackDeep|StackRnn|None:
|
||||
name = project["stack"]["name"]
|
||||
layers = project["stack"]["layers"]
|
||||
|
||||
try:
|
||||
stack_type = StackType[project["stack"]["type_string"]]
|
||||
except KeyError:
|
||||
stack_type = StackType.Deep
|
||||
|
||||
obj = None
|
||||
if stack_type == StackType.Deep:
|
||||
obj = StackDeep(name, work_dir)
|
||||
if stack_type == StackType.Rnn:
|
||||
obj = StackRnn(name, work_dir)
|
||||
|
||||
if obj is None:
|
||||
return obj
|
||||
|
||||
for layer in layers:
|
||||
layer_id = layer["id"]
|
||||
layer_name = layer["name"]
|
||||
num_visible_x = layer["numVisibleX"]
|
||||
num_visible_y = layer["numVisibleY"]
|
||||
num_hidden = layer["numHidden"]
|
||||
try:
|
||||
num_context = layer["numContext"]
|
||||
except KeyError:
|
||||
num_context = 0
|
||||
|
||||
# Determine version by existence of keys
|
||||
params = layer["rbm"]["params"]
|
||||
|
||||
# signal that doSampleBatch will be ignored
|
||||
if params["doSampleBatch"] == 1:
|
||||
raise Exception(f"Error: Parameter \"doSampleBatch\" will be ignored. To continue, set \"doSampleBatch = False\"")
|
||||
|
||||
# import params
|
||||
entity_params = EntityParams.from_dict(params)
|
||||
training_params = TrainingParams.from_dict(params)
|
||||
|
||||
# Create layer
|
||||
layer_obj = Layer(f"{layer_name}-{layer_id}", (num_visible_x, num_visible_y, num_context, num_hidden), entity_params, training_params)
|
||||
|
||||
# Add layer to stack
|
||||
obj.append(layer_obj)
|
||||
|
||||
return obj
|
||||
|
||||
@classmethod
|
||||
def from_file(cls, filename: str, work_dir: str = ".") -> StackDeep|StackRnn:
|
||||
obj = None
|
||||
with open(filename, "r") as fp:
|
||||
prj = json.load(fp)
|
||||
obj = StackFactory.from_dict(prj, work_dir)
|
||||
|
||||
return obj
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
stack = StackFactory.from_file("/home/jens/work/repos/Rbm/test.prj", work_dir="../../results")
|
||||
|
||||
print("Test: [passed]")
|
||||
@@ -0,0 +1,281 @@
|
||||
import numpy as _np_cpu
|
||||
from rbm.status import Status
|
||||
from rbm.train import _to_gpu, cd_binary_binary, cd_gaussian_binary, cd_gaussian_gaussian, cd_binary_gaussian
|
||||
from stack.stack import Stack, StackType
|
||||
from rbm.matrix import Mat, np, rms_error_accu, convert
|
||||
from rbm.entity import Entity, EntityParams, TrainingParams
|
||||
from rbm.layer import Layer
|
||||
|
||||
_CD_FUNC = {
|
||||
Entity.Type.BB_RBM: cd_binary_binary,
|
||||
Entity.Type.GB_RBM: cd_gaussian_binary,
|
||||
Entity.Type.GG_RBM: cd_gaussian_gaussian,
|
||||
Entity.Type.BG_RBM: cd_binary_gaussian,
|
||||
}
|
||||
|
||||
|
||||
class StackRnn(Stack):
|
||||
"""Recurrent RBM stack.
|
||||
|
||||
Two modes selected by the number of appended layers:
|
||||
|
||||
Shared weights (1 layer):
|
||||
The same Entity processes every time step.
|
||||
Sequences may have any length.
|
||||
|
||||
Unrolled / own weights (N layers):
|
||||
layers[t % N] processes time step t — each position has its own W.
|
||||
Training requires sequences of exactly length N.
|
||||
"""
|
||||
|
||||
def __init__(self, name: str, work_dir: str = '.'):
|
||||
Stack.__init__(self, StackType.Rnn, name, work_dir)
|
||||
self._h: Mat | None = None # context / hidden state: (batch, h_size)
|
||||
self._t: int = 0 # current time-step counter
|
||||
|
||||
# ── Convenience factories ─────────────────────────────────────────────
|
||||
|
||||
@staticmethod
|
||||
def make_layer(name: str, sensory_size: int, h_size: int,
|
||||
entity_params: EntityParams, training_params: TrainingParams) -> Layer:
|
||||
"""Single layer for shared-weights mode."""
|
||||
return Layer(name, (1, sensory_size, h_size, h_size), entity_params, training_params)
|
||||
|
||||
@staticmethod
|
||||
def make_unrolled(time_steps: int, sensory_size: int, h_size: int,
|
||||
entity_params: EntityParams, training_params: TrainingParams) -> list[Layer]:
|
||||
"""N layers for unrolled (own-weights) mode — one per time step.
|
||||
|
||||
Usage::
|
||||
|
||||
for layer in StackRnn.make_unrolled(T, ...):
|
||||
rnn.append(layer)
|
||||
"""
|
||||
return [
|
||||
Layer(f"t{t}", (1, sensory_size, h_size, h_size), entity_params, training_params)
|
||||
for t in range(time_steps)
|
||||
]
|
||||
|
||||
# ── Mode ─────────────────────────────────────────────────────────────
|
||||
|
||||
@property
|
||||
def is_shared(self) -> bool:
|
||||
return self.num_layers() == 1
|
||||
|
||||
def _entity_at(self, t: int) -> Entity:
|
||||
return self.from_index(t % self.num_layers()).entity
|
||||
|
||||
def next_entity(self) -> Entity:
|
||||
"""Entity that will be used by the next step() call."""
|
||||
return self._entity_at(self._t)
|
||||
|
||||
def current_entity(self) -> Entity:
|
||||
"""Entity used by the most recent step() call."""
|
||||
return self._entity_at(max(0, self._t - 1))
|
||||
|
||||
# ── Derived sizes ─────────────────────────────────────────────────────
|
||||
|
||||
def h_size(self, layer_idx: int = 0) -> int:
|
||||
return self.from_index(layer_idx).entity.shape[1]
|
||||
|
||||
def sensory_size(self, layer_idx: int = 0) -> int:
|
||||
e = self.from_index(layer_idx).entity
|
||||
return e.shape[0] - e.shape[1]
|
||||
|
||||
# ── State management ──────────────────────────────────────────────────
|
||||
|
||||
def reset(self, batch_size: int = 1):
|
||||
"""Zero the context vector and reset the time-step counter."""
|
||||
self._h = np.zeros((batch_size, self.h_size()))
|
||||
self._t = 0
|
||||
|
||||
# ── Inference ─────────────────────────────────────────────────────────
|
||||
|
||||
def step(self, x: Mat) -> Mat:
|
||||
"""One time step forward.
|
||||
|
||||
In shared mode uses the single entity.
|
||||
In unrolled mode uses layers[_t % N] and advances _t.
|
||||
|
||||
x: (batch_size, sensory_size) or (sensory_size,)
|
||||
Returns the new context vector h_t.
|
||||
"""
|
||||
if x.ndim == 1:
|
||||
x = x[None, :]
|
||||
batch_size = x.shape[0]
|
||||
if self._h is None or self._h.shape[0] != batch_size:
|
||||
self.reset(batch_size)
|
||||
|
||||
entity = self._entity_at(self._t)
|
||||
visible = np.concatenate([self._h, _to_gpu(x)], axis=1)
|
||||
self._h = entity.forward(visible)
|
||||
self._t += 1
|
||||
return self._h
|
||||
|
||||
def reconstruct(self, h: Mat) -> Mat:
|
||||
"""Decode h → visible, returning only the sensory portion.
|
||||
|
||||
Uses the entity from the most recent step() call.
|
||||
"""
|
||||
entity = self.current_entity()
|
||||
visible = entity.reconstruct(h)
|
||||
return visible[:, entity.shape[1]:]
|
||||
|
||||
# ── Training ──────────────────────────────────────────────────────────
|
||||
|
||||
def train(self, sequences: Mat, status: Status = None):
|
||||
"""Train on sequences.
|
||||
|
||||
sequences: (T, sensory_size) — single sequence
|
||||
(num_seq, T, sensory_size) — batch of sequences
|
||||
|
||||
Shared mode (1 layer): T may be any value.
|
||||
Unrolled mode (N layers): T must equal N.
|
||||
"""
|
||||
if status is None:
|
||||
status = Status()
|
||||
|
||||
seqs = sequences if sequences.ndim == 3 else sequences[None, :]
|
||||
num_seq, T, _ = seqs.shape
|
||||
|
||||
if self.is_shared:
|
||||
entity = self.from_index(0).entity
|
||||
print(f"Train shared ({entity.name}) "
|
||||
f"for {entity.training_params.num_epochs} epochs")
|
||||
if entity.enable_training and entity.training_params is not None:
|
||||
self._train_shared(entity, seqs, num_seq, T, status)
|
||||
else:
|
||||
assert T == self.num_layers(), (
|
||||
f"Unrolled mode: sequence length T={T} "
|
||||
f"must equal num_layers={self.num_layers()}"
|
||||
)
|
||||
params = self.from_index(0).entity.training_params
|
||||
print(f"Train unrolled ({self.num_layers()} layers) "
|
||||
f"for {params.num_epochs} epochs")
|
||||
self._train_unrolled(seqs, num_seq, T, status)
|
||||
|
||||
def _train_shared(self, entity: Entity, seqs: Mat,
|
||||
num_seq: int, T: int, status: Status):
|
||||
"""One entity, reused at every time step."""
|
||||
cd_func = _CD_FUNC[entity.type]
|
||||
params = entity.training_params
|
||||
h_sz = entity.shape[1]
|
||||
|
||||
d_progress = 100.0 / params.num_epochs
|
||||
progress = 0.0
|
||||
keep_running = True
|
||||
|
||||
entity.grad_zero()
|
||||
status.on_change(entity)
|
||||
|
||||
for epoch in range(params.num_epochs):
|
||||
if not keep_running:
|
||||
break
|
||||
|
||||
h = np.zeros((num_seq, h_sz))
|
||||
err_total = 0.0
|
||||
|
||||
for t in range(T):
|
||||
x_t = _to_gpu(seqs[:, t, :])
|
||||
visible = np.concatenate([h, x_t], axis=1)
|
||||
|
||||
dwhv, dbv, dbh = cd_func(entity, visible)
|
||||
grad = entity.grad_compute(dbv, dbh, dwhv)
|
||||
entity.state_adjust(grad, 1.0 / num_seq)
|
||||
|
||||
h = entity.forward(visible)
|
||||
err_total += rms_error_accu(visible - entity.reconstruct(h))
|
||||
|
||||
progress += d_progress
|
||||
if status.want_report(round(progress)):
|
||||
if not status.on_change(entity, {
|
||||
"progress": {"value": round(progress), "unit": "%"},
|
||||
"err_rms": {"value": err_total / T, "unit": ""},
|
||||
}):
|
||||
keep_running = False
|
||||
break
|
||||
|
||||
h = np.zeros((num_seq, h_sz))
|
||||
err_total = 0.0
|
||||
for t in range(T):
|
||||
x_t = _to_gpu(seqs[:, t, :])
|
||||
visible = np.concatenate([h, x_t], axis=1)
|
||||
h = entity.forward(visible)
|
||||
err_total += rms_error_accu(visible - entity.reconstruct(h))
|
||||
status.on_change(entity, {
|
||||
"progress": {"value": 100, "unit": "%"},
|
||||
"err_rms_total": {"value": err_total / T, "unit": ""},
|
||||
})
|
||||
|
||||
def _train_unrolled(self, seqs: Mat, num_seq: int, T: int, status: Status):
|
||||
"""N entities, one per time step — each has its own W, b_v, b_h.
|
||||
|
||||
Uses temporal-shift padding (matching C++ RnnStack):
|
||||
Flatten (num_seq, T, sensory_size) → (N, sensory_size), append T-1 zero
|
||||
rows, then layer t trains on batch_padded[t : N+t] — a one-step delay.
|
||||
|
||||
Joint training: all layers are updated together each epoch.
|
||||
Context c from layer t feeds layer t+1 within the same epoch pass,
|
||||
so gradients propagate through the full temporal chain.
|
||||
"""
|
||||
params = self.from_index(0).entity.training_params
|
||||
h_sz = self.from_index(0).entity.shape[1]
|
||||
s_sz = seqs.shape[2]
|
||||
d_progress = 100.0 / params.num_epochs
|
||||
progress = 0.0
|
||||
keep_running = True
|
||||
|
||||
# Flatten to (N, s_sz) keeping sequence-major order, on CPU for slicing
|
||||
flat = _np_cpu.asarray(convert(seqs) if hasattr(seqs, 'get') else seqs)
|
||||
flat = flat.reshape(-1, s_sz)
|
||||
N = flat.shape[0]
|
||||
pad = _np_cpu.zeros((T - 1, s_sz), dtype=flat.dtype)
|
||||
batch_pad = _np_cpu.concatenate([flat, pad], axis=0) # (N+T-1, s_sz)
|
||||
|
||||
for layer in self.layers:
|
||||
layer.entity.grad_zero()
|
||||
status.on_change(self.from_index(0).entity)
|
||||
|
||||
for epoch in range(params.num_epochs):
|
||||
if not keep_running:
|
||||
break
|
||||
|
||||
c = np.zeros((N, h_sz))
|
||||
err_total = 0.0
|
||||
|
||||
for t, layer in enumerate(self.layers):
|
||||
entity = layer.entity
|
||||
cd_func = _CD_FUNC[entity.type]
|
||||
|
||||
x_t = _to_gpu(batch_pad[t : N + t]) # shifted slice, (N, s_sz)
|
||||
visible = np.concatenate([c, x_t], axis=1)
|
||||
|
||||
dwhv, dbv, dbh = cd_func(entity, visible)
|
||||
grad = entity.grad_compute(dbv, dbh, dwhv)
|
||||
entity.state_adjust(grad, 1.0 / N)
|
||||
|
||||
c = entity.forward(visible)
|
||||
err_total += rms_error_accu(visible - entity.reconstruct(c))
|
||||
|
||||
progress += d_progress
|
||||
if status.want_report(round(progress)):
|
||||
if not status.on_change(self.from_index(0).entity, {
|
||||
"progress": {"value": round(progress), "unit": "%"},
|
||||
"err_rms": {"value": err_total / T, "unit": ""},
|
||||
}):
|
||||
keep_running = False
|
||||
break
|
||||
|
||||
# Final report pass
|
||||
c = np.zeros((N, h_sz))
|
||||
err_total = 0.0
|
||||
for t, layer in enumerate(self.layers):
|
||||
entity = layer.entity
|
||||
x_t = _to_gpu(batch_pad[t : N + t])
|
||||
visible = np.concatenate([c, x_t], axis=1)
|
||||
c = entity.forward(visible)
|
||||
err_total += rms_error_accu(visible - entity.reconstruct(c))
|
||||
status.on_change(self.from_index(0).entity, {
|
||||
"progress": {"value": 100, "unit": "%"},
|
||||
"err_rms_total": {"value": err_total / T, "unit": ""},
|
||||
})
|
||||
@@ -0,0 +1,40 @@
|
||||
import numpy as np
|
||||
|
||||
VOCAB = ' .!?ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789'
|
||||
_CH2IDX = {ch: i for i, ch in enumerate(VOCAB)}
|
||||
|
||||
def ch2idx(ch: str) -> int:
|
||||
return _CH2IDX[ch]
|
||||
|
||||
def idx2ch(idx: int) -> str:
|
||||
return VOCAB[idx]
|
||||
|
||||
|
||||
def concat(v: np.ndarray, c: np.ndarray) -> np.ndarray:
|
||||
return np.concatenate([v, c])
|
||||
|
||||
def split(vc: np.ndarray, v_size: np.ndarray|int) -> tuple[np.ndarray, np.ndarray]:
|
||||
n = len(v_size) if isinstance(v_size, np.ndarray) else v_size
|
||||
return vc[:n], vc[n:]
|
||||
|
||||
def shift_right(m: np.ndarray) -> np.ndarray:
|
||||
return np.hstack([np.zeros((m.shape[0], 1)), m[:, :-1]])
|
||||
|
||||
def shift_left(m: np.ndarray) -> np.ndarray:
|
||||
return np.hstack([m[:, 1:], np.zeros((m.shape[0], 1))])
|
||||
|
||||
def shift_up(m: np.ndarray) -> np.ndarray:
|
||||
return np.vstack([m[1:, :], np.zeros((1, m.shape[1]))])
|
||||
|
||||
def shift_down(m: np.ndarray) -> np.ndarray:
|
||||
return np.vstack([np.zeros((1, m.shape[1])), m[:-1, :]])
|
||||
|
||||
def clamp_one_hot(src_dst: np.ndarray) -> np.ndarray:
|
||||
index = np.argmax(src_dst)
|
||||
result = np.zeros_like(src_dst)
|
||||
result[index] = 1
|
||||
return result
|
||||
|
||||
class Rnn:
|
||||
def __init__(self):
|
||||
pass
|
||||
@@ -0,0 +1,56 @@
|
||||
import os
|
||||
from enum import Enum
|
||||
from rbm.layer import Layer
|
||||
|
||||
class StackType(Enum):
|
||||
Deep = "Deep",
|
||||
Rnn = "Rnn"
|
||||
|
||||
class StackException(Exception):
|
||||
pass
|
||||
|
||||
class Stack:
|
||||
def __init__(self, stack_type: StackType, name: str, work_dir: str = '.'):
|
||||
self.stack_type = stack_type
|
||||
self.name = name
|
||||
self.work_dir = work_dir
|
||||
self.layers: list[Layer] = []
|
||||
os.makedirs(self.work_dir, exist_ok=True)
|
||||
|
||||
def num_layers(self):
|
||||
return len(self.layers)
|
||||
|
||||
def append(self, layer: Layer):
|
||||
self.layers.append(layer)
|
||||
return self.num_layers()-1
|
||||
|
||||
def remove(self, layer: Layer):
|
||||
for index, lay in enumerate(self.layers):
|
||||
if lay == layer:
|
||||
del(self.layers[index])
|
||||
return
|
||||
raise StackException(f"Layer \"{layer.name}\" does not exist")
|
||||
|
||||
def from_index(self, index : int) -> Layer:
|
||||
result = self.layers[index]
|
||||
return result
|
||||
|
||||
def from_name(self, name: str):
|
||||
for layer in self.layers:
|
||||
if layer.name in name:
|
||||
return layer
|
||||
return None
|
||||
|
||||
def state_init(self, std: float):
|
||||
for layer in self.layers:
|
||||
layer.init(std)
|
||||
|
||||
def state_save(self):
|
||||
for index, layer in enumerate(self.layers):
|
||||
filepath = os.path.join(self.work_dir, f"{self.name}-{index}-state.npz")
|
||||
layer.save(filepath)
|
||||
|
||||
def state_load(self):
|
||||
for index, layer in enumerate(self.layers):
|
||||
filepath = os.path.join(self.work_dir, f"{self.name}-{index}-state.npz")
|
||||
layer.load(filepath)
|
||||
Reference in New Issue
Block a user