Even the later "fat processor" (RISC) Connection Machines had several hundreds of CPU (the earlier models had tens of thousands of single-bit processors)
Today most software can only scale if you add a faster processor, that has to stop, before that happens, there is not much the processor manufacturer can do, their hands are tied IMO
The programming interface of a modern x86 CPU is best thought of as a virtual machine with a JIT. The JIT performs some optimizations to make old code run more in parallel without the need for recompilation. This is a terrible model for compilers and programmers, since it makes optimizations hit or miss, and of course it's why we're in this whole spectre/meltdown mess at the moment. If we switched to a more reasonable programming model which actually met the needs of software (fine grained interprocessor communication without going through central memory, many more registers, actual software access to pipelines, caches, etc.) then it would make a lot of old software run slower (through emulation), but new software could finally make better use of your hardware...
Seeing how Itanium went, their risk aversion seems pretty wise.
Itanium seems to satisfy most of your desires with its VLIW/EPIC architecture, which exposes much more to the compiler.
> Today most software can only scale if you add a faster processor, that has to stop, before that happens, there is not much the processor manufacturer can do, their hands are tied IMO
These are not recent developments, chip makers could have opted not to use them and software developers would have done otherwise.
The hardware is developed to make the existing software faster, but the software is also developed to run well on the current hardware.
We have some idea of how to do massively-parallel programming. That is, after all, what a supercomputer is. But one of the things we've found is that communication doesn't scale. For very low core counts, you can hook up each core to talk to every other one at the same latency--basically, O(N)-degree topology. For slightly higher core counts, you can keep it a hypercube, which is O(lg N) degree. But when you start thinking about a few hundred cores, you end up with a mesh, which is O(1)-degree.
What that means is you struggle to scale problems that have very high communication costs compared to computation. Matrix multiplication is wonderful--you're doing O(N^1.5) computations for each data, so communication goes down as size goes up [1]. But a BFS graph traversal is horrible, because your communication costs go up if average degree is more than O(1).
[1] This is why the TOP500 benchmark is to some degree bullshit. It's measuring performance on applications that are fundamentally computation-bound, not communication bound, so it tends to penalize machines that focus on improving communication bandwidth, which is often helpful for many HPC applications.